diff --git a/README.md b/README.md new file mode 100644 index 0000000..379a498 --- /dev/null +++ b/README.md @@ -0,0 +1,67 @@ +# IoT Consumer — обработчик событий + +Читает сообщения из очереди RabbitMQ и сохраняет их в Redis (оперативный кэш) и MongoDB (постоянный архив). + +## Пайплайн + +``` +Producer → RabbitMQ → Consumer (этот сервис) → Redis + MongoDB → Dashboard +``` + +## Хранилища + +### Redis — оперативный кэш + +| Ключ | Тип | Содержимое | +|---|---|---| +| `iot:latest` | Hash | Последнее значение каждого датчика (`sensor_id → JSON`) | +| `iot:counters` | Hash | Счётчики событий по типам (`temperature: N`, ...) | +| `iot:recent` | ZSET | Лента последних 1000 событий (сортировка по времени) | + +### MongoDB — постоянный архив + +- Коллекция: `iot_events` +- TTL-индекс: автоматическое удаление через **7 дней** + +## Переменные окружения + +| Переменная | По умолчанию | Описание | +|---|---|---| +| `RMQ_HOST` | `localhost` | Хост RabbitMQ | +| `RMQ_PORT` | `5672` | Порт AMQP | +| `RMQ_USER` | `guest` | Логин RabbitMQ | +| `RMQ_PASS` | `guest` | Пароль RabbitMQ | +| `RMQ_QUEUE` | `iot-events` | Имя очереди | +| `REDIS_HOST` | `localhost` | Хост Redis | +| `REDIS_PORT` | `6379` | Порт Redis | +| `REDIS_PASS` | — | Пароль Redis (необязательно) | +| `MONGO_URI` | `mongodb://localhost:27017/iot` | Строка подключения MongoDB | +| `MONGO_DB` | `iot` | Имя базы MongoDB | +| `PORT` | `3000` | Порт HTTP | + +## Эндпоинты + +| Путь | Ответ | +|---|---| +| `GET /` | `{"service":"iot-consumer","consumed":N,"errors":N}` | +| `GET /health` | `{"status":"ok"}` | + +## Локальный запуск + +```bash +npm install +# Нужны RabbitMQ, Redis и MongoDB (локально или в Docker) +RMQ_HOST=localhost REDIS_HOST=localhost MONGO_URI=mongodb://localhost:27017/iot node server.js +``` + +## Деплой на Nubes Cloud + +1. Запушьте этот репо в свой Gitea +2. В Terraform-конфиге укажите `git_path` на ваш репо и `git_revision` +3. Nubes автоматически запустит `npm start` + +## Отказоустойчивость + +- Каждое хранилище (Redis, MongoDB) обёрнуто в отдельный `try/catch` +- Отказ одного хранилища **не ломает** запись в другое +- Ошибки логируются, consumer продолжает работу diff --git a/package.json b/package.json index 7366070..3ceb00b 100644 --- a/package.json +++ b/package.json @@ -2,6 +2,13 @@ "name": "iot-consumer", "version": "1.0.0", "main": "server.js", - "scripts": { "start": "node server.js" }, - "dependencies": { "amqplib": "^0.10.0", "express": "^4.21.0", "redis": "^4.7.0", "mongodb": "^6.0.0" } + "scripts": { + "start": "node server.js" + }, + "dependencies": { + "amqplib": "^0.10.0", + "express": "^4.21.0", + "redis": "^4.7.0", + "mongodb": "^6.0.0" + } } diff --git a/requirements.txt b/requirements.txt deleted file mode 100644 index 9cc1ecd..0000000 --- a/requirements.txt +++ /dev/null @@ -1,4 +0,0 @@ -Flask>=3.0 -gunicorn>=21.2 -redis>=5.0 -pymongo>=4.0 diff --git a/server.js b/server.js index 91c751a..ff72062 100644 --- a/server.js +++ b/server.js @@ -1,27 +1,92 @@ -const amqp = require("amqplib"); -const express = require("express"); -const { createClient } = require("redis"); -const { MongoClient } = require("mongodb"); +// ============================================================================= +// IoT Consumer — Обработчик событий из RabbitMQ +// ============================================================================= +// Читает сообщения из очереди RabbitMQ (iot-events), помещённые туда +// Producer'ом, и сохраняет их в два хранилища: +// 1. Redis — оперативный кэш для Dashboard: +// - iot:latest (hash) — последнее значение каждого датчика +// - iot:counters (hash) — счётчики событий по типам устройств +// - iot:recent (zset) — лента последних 1000 событий +// 2. MongoDB — постоянный архив всех событий: +// - коллекция iot_events с TTL-индексом (автоудаление через 7 дней) +// +// Каждое хранилище обёрнуто в try/catch — отказ одного не ломает другое. +// +// Переменные окружения (задаются Terraform'ом через json_env): +// RMQ_HOST — хост RabbitMQ +// RMQ_PORT — порт AMQP (5672) +// RMQ_USER — логин RabbitMQ +// RMQ_PASS — пароль RabbitMQ +// RMQ_QUEUE — имя очереди (iot-events) +// REDIS_HOST — хост Redis +// REDIS_PORT — порт Redis (6379) +// REDIS_PASS — пароль Redis +// MONGO_URI — строка подключения MongoDB (mongodb://user:pass@host:port/db) +// MONGO_DB — имя базы MongoDB (iot) +// PORT — порт HTTP-сервера (3000) +// +// Эндпоинты: +// GET / — статистика (consumed, errors) +// GET /health — health-check для Kubernetes +// ============================================================================= -const RMQ_HOST = process.env.RMQ_HOST || "localhost"; -const RMQ_PORT = parseInt(process.env.RMQ_PORT || "5672"); -const RMQ_USER = process.env.RMQ_USER || "guest"; -const RMQ_PASS = process.env.RMQ_PASS || "guest"; -const RMQ_QUEUE = process.env.RMQ_QUEUE || "iot-events"; +// --------------------------------------------------------------------------- +// Зависимости +// --------------------------------------------------------------------------- +const amqp = require("amqplib"); // AMQP-клиент для RabbitMQ +const express = require("express"); // HTTP-фреймворк для /health и /stats +const { createClient } = require("redis"); // Redis-клиент (v4+, с Promise API) +const { MongoClient } = require("mongodb");// MongoDB-клиент (v6+, нативный драйвер) -const REDIS_HOST = process.env.REDIS_HOST || "localhost"; -const REDIS_PORT = parseInt(process.env.REDIS_PORT || "6379"); -const REDIS_PASS = process.env.REDIS_PASS || undefined; +// --------------------------------------------------------------------------- +// Конфигурация RabbitMQ из переменных окружения +// --------------------------------------------------------------------------- +const RMQ_HOST = process.env.RMQ_HOST || "localhost"; // K8s-хост +const RMQ_PORT = parseInt(process.env.RMQ_PORT || "5672"); // AMQP-порт +const RMQ_USER = process.env.RMQ_USER || "guest"; // логин +const RMQ_PASS = process.env.RMQ_PASS || "guest"; // пароль +const RMQ_QUEUE = process.env.RMQ_QUEUE || "iot-events"; // имя очереди +// --------------------------------------------------------------------------- +// Конфигурация Redis из переменных окружения +// --------------------------------------------------------------------------- +const REDIS_HOST = process.env.REDIS_HOST || "localhost"; // K8s-хост +const REDIS_PORT = parseInt(process.env.REDIS_PORT || "6379"); // порт Redis +const REDIS_PASS = process.env.REDIS_PASS || undefined; // пароль (undefined = без пароля) + +// --------------------------------------------------------------------------- +// Конфигурация MongoDB из переменных окружения +// --------------------------------------------------------------------------- +// MONGO_URI формируется Terraform'ом как: +// mongodb://admin:@хост:27017/iot?authSource=admin +// (без пароля для демо-стенда — platform limitation) +// --------------------------------------------------------------------------- const MONGO_URI = process.env.MONGO_URI || "mongodb://localhost:27017/iot"; -const MONGO_DB = process.env.MONGO_DB || "iot"; +const MONGO_DB = process.env.MONGO_DB || "iot"; -let consumed = 0, errors = 0; -let redis, mongo, mongoColl; +// --------------------------------------------------------------------------- +// Глобальное состояние +// --------------------------------------------------------------------------- +let consumed = 0, errors = 0; // счётчики для / +let redis, // Redis-клиент (ленивая инициализация) + mongo, // MongoDB-клиент (ленивая инициализация) + mongoColl; // MongoDB-коллекция iot_events +// --------------------------------------------------------------------------- +// getRedis() — ленивое подключение к Redis +// --------------------------------------------------------------------------- +// Создаёт клиент Redis один раз, при первом вызове. +// Подавляем ошибки Redis (redis.on("error", () => {})) — не хотим крашить +// consumer при временной недоступности Redis. +// --------------------------------------------------------------------------- async function getRedis() { if (!redis) { - redis = createClient({ socket: { host: REDIS_HOST, port: REDIS_PORT }, password: REDIS_PASS }); + // Создаём клиент с настройками из переменных окружения + redis = createClient({ + socket: { host: REDIS_HOST, port: REDIS_PORT }, // TCP-сокет + password: REDIS_PASS, // пароль (undefined если нет) + }); + // Подавляем ошибки Redis — не крашим consumer из-за временных проблем redis.on("error", () => {}); await redis.connect(); console.log("Redis connected"); @@ -29,61 +94,175 @@ async function getRedis() { return redis; } +// --------------------------------------------------------------------------- +// getMongo() — ленивое подключение к MongoDB +// --------------------------------------------------------------------------- +// Создаёт клиент MongoDB один раз, при первом вызове. +// При подключении создаёт TTL-индекс на поле timestamp: +// expireAfterSeconds: 604800 = 7 дней +// Это значит MongoDB автоматически удалит документы старше 7 дней — +// не нужно ручной очистки. +// --------------------------------------------------------------------------- async function getMongo() { if (!mongo) { - mongo = new MongoClient(MONGO_URI); + mongo = new MongoClient(MONGO_URI); // строка подключения из env await mongo.connect(); + // Получаем коллекцию iot_events в базе MONGO_DB mongoColl = mongo.db(MONGO_DB).collection("iot_events"); + // TTL-индекс: документы старше 7 дней удаляются автоматически await mongoColl.createIndex({ timestamp: 1 }, { expireAfterSeconds: 604800 }); console.log("MongoDB connected"); } return mongoColl; } +// --------------------------------------------------------------------------- +// store(event) — сохраняет событие в Redis и MongoDB +// --------------------------------------------------------------------------- +// Каждое хранилище — в своём try/catch: +// - Если Redis упал — событие всё равно попадёт в MongoDB +// - Если MongoDB упала — событие всё равно попадёт в Redis +// Это гарантирует, что отказ одного хранилища не блокирует другое. +// +// Redis-структуры: +// iot:latest — HSET sensor_id → JSON{value, unit, location, timestamp} +// (каждый новый запрос перезаписывает предыдущий) +// iot:counters — HINCRBY device_type → +1 +// (счётчик событий по типам: temperature, humidity, power_meter) +// iot:recent — ZADD score=Date.now() → JSON-событие +// ZREMRANGEBYRANK 0..-1001 — держим только 1000 последних +// --------------------------------------------------------------------------- async function store(event) { + // === Запись в Redis (оперативный кэш) === try { const r = await getRedis(); - await r.hSet("iot:latest", event.sensor_id, JSON.stringify({ - value: event.value, unit: event.unit, location: event.location, timestamp: event.timestamp, - })); - await r.hIncrBy("iot:counters", event.device_type, 1); - await r.zAdd("iot:recent", [{ score: Date.now(), value: JSON.stringify(event) }]); - await r.zRemRangeByRank("iot:recent", 0, -1001); - } catch (e) { console.error("Redis error:", e.message); } + // HSET — последнее значение датчика (ключ = sensor_id) + // Перезаписывает предыдущее значение — всегда актуальное + await r.hSet("iot:latest", event.sensor_id, JSON.stringify({ + value: event.value, + unit: event.unit, + location: event.location, + timestamp: event.timestamp, + })); + + // HINCRBY — атомарный инкремент счётчика по типу устройства + await r.hIncrBy("iot:counters", event.device_type, 1); + + // ZADD — добавляем в сортированное множество (счёт = timestamp) + // score = Date.now() для сортировки по времени + await r.zAdd("iot:recent", [{ + score: Date.now(), + value: JSON.stringify(event), + }]); + + // ZREMRANGEBYRANK — удаляем всё что выходит за пределы последних 1000 + // 0..-1001 означает: удалить с 0-го элемента до элемента с индексом + // (len - 1001), оставляя 1000 последних + await r.zRemRangeByRank("iot:recent", 0, -1001); + } catch (e) { + // Redis недоступен — логируем, но НЕ крашим consumer + console.error("Redis error:", e.message); + } + + // === Запись в MongoDB (постоянный архив) === try { const coll = await getMongo(); + + // insertOne — вставка одного документа + // timestamp преобразуется из ISO-строки в Date для TTL-индекса await coll.insertOne({ - event_id: event.event_id, sensor_id: event.sensor_id, - device_type: event.device_type, location: event.location, - value: event.value, unit: event.unit, timestamp: new Date(event.timestamp), + event_id: event.event_id, + sensor_id: event.sensor_id, + device_type: event.device_type, + location: event.location, + value: event.value, + unit: event.unit, + timestamp: new Date(event.timestamp), // Date-объект для TTL }); - } catch (e) { console.error("Mongo error:", e.message); } + } catch (e) { + // MongoDB недоступна — логируем, но НЕ крашим consumer + console.error("Mongo error:", e.message); + } } +// --------------------------------------------------------------------------- +// start() — подключение к RabbitMQ и запуск консьюмера +// --------------------------------------------------------------------------- +// 1. Подключается к RabbitMQ по AMQP +// 2. Создаёт канал и гарантирует существование очереди +// 3. Устанавливает prefetch(10) — не более 10 неподтверждённых сообщений +// одновременно (backpressure: не загружаем consumer слишком сильно) +// 4. consume() — подписывается на очередь, обрабатывает каждое сообщение: +// a. Парсит JSON из тела сообщения +// b. Вызывает store() для записи в Redis + MongoDB +// c. ack() — подтверждает обработку (сообщение удаляется из очереди) +// d. При ошибке: nack(msg, false, true) — возвращает в очередь (requeue) +// --------------------------------------------------------------------------- async function start() { - const conn = await amqp.connect(`amqp://${RMQ_USER}:${RMQ_PASS}@${RMQ_HOST}:${RMQ_PORT}`); + // Подключение к RabbitMQ + const conn = await amqp.connect( + `amqp://${RMQ_USER}:${RMQ_PASS}@${RMQ_HOST}:${RMQ_PORT}` + ); const ch = await conn.createChannel(); + + // Гарантируем существование очереди (идемпотентно) await ch.assertQueue(RMQ_QUEUE, { durable: true }); + + // prefetch(10): не отправлять больше 10 сообщений одновременно, + // пока предыдущие не подтверждены (ack) + // Это защита от перегрузки consumer'а ch.prefetch(10); + + // Подписка на очередь await ch.consume(RMQ_QUEUE, async (msg) => { + // msg === null означает отмену подписки (например, очередь удалена) if (!msg) return; + try { + // Парсим JSON из тела сообщения (Buffer → string → object) const event = JSON.parse(msg.content.toString()); + + // Сохраняем в Redis + MongoDB await store(event); - consumed++; - ch.ack(msg); + + consumed++; // инкрементируем счётчик успешных + ch.ack(msg); // подтверждаем — сообщение удаляется из очереди + + // Каждые 10 сообщений — логируем прогресс if (consumed % 10 === 0) console.log(`Consumed ${consumed} events`); } catch (e) { errors++; + // nack с requeue=true — возвращаем сообщение в очередь + // для повторной попытки (если ошибка временная) + // ВНИМАНИЕ: при систематической ошибке может вызвать infinite loop! ch.nack(msg, false, true); console.error("Error:", e.message); } }); + console.log(`Consumer listening on queue: ${RMQ_QUEUE}`); } +// --------------------------------------------------------------------------- +// HTTP-сервер (Express) +// --------------------------------------------------------------------------- const app = express(); -app.get("/", (req, res) => res.json({ service: "iot-consumer", status: "running", queue: RMQ_QUEUE, consumed, errors })); + +// GET / — статистика consumer'а +app.get("/", (req, res) => res.json({ + service: "iot-consumer", + status: "running", + queue: RMQ_QUEUE, + consumed, + errors, +})); + +// GET /health — для Kubernetes liveness/readiness probe app.get("/health", (req, res) => res.json({ status: "ok" })); -app.listen(process.env.PORT || 3000, () => { console.log("Consumer started"); start(); }); + +// Запуск: слушаем порт, затем подключаемся к RabbitMQ +app.listen(process.env.PORT || 3000, () => { + console.log("Consumer started"); + start(); // подключаемся к RabbitMQ и начинаем обрабатывать сообщения +});