// ============================================================================= // 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 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+, нативный драйвер) // --------------------------------------------------------------------------- // Конфигурация 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"; // --------------------------------------------------------------------------- // Глобальное состояние // --------------------------------------------------------------------------- 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 }, // TCP-сокет password: REDIS_PASS, // пароль (undefined если нет) }); // Подавляем ошибки Redis — не крашим consumer из-за временных проблем redis.on("error", () => {}); await redis.connect(); console.log("Redis connected"); } return redis; } // --------------------------------------------------------------------------- // getMongo() — ленивое подключение к MongoDB // --------------------------------------------------------------------------- // Создаёт клиент MongoDB один раз, при первом вызове. // При подключении создаёт TTL-индекс на поле timestamp: // expireAfterSeconds: 604800 = 7 дней // Это значит MongoDB автоматически удалит документы старше 7 дней — // не нужно ручной очистки. // --------------------------------------------------------------------------- async function getMongo() { if (!mongo) { 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(); // 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), // Date-объект для TTL }); } 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() { // Подключение к 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); // подтверждаем — сообщение удаляется из очереди // Каждые 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(); // 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" })); // Запуск: слушаем порт, затем подключаемся к RabbitMQ app.listen(process.env.PORT || 3000, () => { console.log("Consumer started"); start(); // подключаемся к RabbitMQ и начинаем обрабатывать сообщения });