Стримінгова AI-система детекції аномалій IoT-датчиків

Проектуємо та впроваджуємо системи штучного інтелекту: від прототипу до production-ready рішення. Наша команда поєднує експертизу в машинному навчанні, дата-інжинірингу та MLOps, щоб AI працював не в лабораторії, а в реальному бізнесі.
Показано 1 з 1Усі 1564 послуг
Стримінгова AI-система детекції аномалій IoT-датчиків
Середній
~2-4 тижні
Часті запитання

Напрямки AI-розробки

Етапи розробки AI-рішення

Останні роботи

  • image_website-b2b-advance_0.webp
    Розробка сайту компанії B2B ADVANCE
    1357
  • image_web-applications_feedme_466_0.webp
    Розробка веб-додатків для компанії FEEDME
    1250
  • image_websites_belfingroup_462_0.webp
    Розробка веб-сайту для компанії БЕЛФІНГРУП
    956
  • image_ecommerce_furnoro_435_0.webp
    Розробка інтернет магазину для компанії FURNORO
    1188
  • image_logo-advance_0.webp
    Розробка логотипу компанії B2B Advance
    646
  • image_crm_enviok_479_0.webp
    Розробка веб-додатків для компанії Enviok
    929

На виробничій лінії датчик температури раптово показує 95°C замість робочих 60°C. Це фізична аномалія (перегрів підшипника) чи збій сенсора? Кожне хибне спрацювання знижує довіру операторів, а пропущена аварія коштує мільйони. Ми створили стримінгову систему на Kafka, Flink та ONNX, яка вирішує цю дилему з точністю 95% і знижує хибні тривоги в 3 рази. Наш досвід — понад 50 впроваджень у промисловості та енергетиці.

Традиційні порогові методи дають до 40% хибних тривог. Контекстуальні аномалії, що враховують час доби та день тижня, знижують цей показник до 5%. Ми комбінуємо статистику (z-score, EWMA) та ML-інференс на edge-пристроях, щоб мінімізувати затримку та зберегти точність на рівні 95%. Типовий стек: MQTT-брокер Mosquitto, Kafka 3.5 для буферизації, Flink для віконної обробки, PyTorch для навчання та ONNX Runtime для інференсу на edge. Інтеграція з Grafana + InfluxDB для моніторингу в реальному часі. Згідно з документацією Apache Kafka, така архітектура гарантує відмовостійкість та масштабованість.

Як працює стримінгова архітектура?

Pipeline від датчика до алерту:

MQTT (датчик) → Kafka → Flink / Spark Streaming → ML inference → AlertManager
                              ↓
                         InfluxDB / TimescaleDB
                              ↓
                         Grafana Dashboard

Схема обробки Kafka:

from kafka import KafkaConsumer, KafkaProducer
import json
import numpy as np
from collections import defaultdict, deque

class IoTAnomalyProcessor:
    def __init__(self, bootstrap_servers='kafka:9092',
                 window_size=60):  # 60 останніх значень
        self.consumer = KafkaConsumer(
            'iot-sensor-raw',
            bootstrap_servers=bootstrap_servers,
            value_deserializer=lambda m: json.loads(m.decode()),
            group_id='anomaly-detection'
        )
        self.producer = KafkaProducer(
            bootstrap_servers=bootstrap_servers,
            value_serializer=lambda v: json.dumps(v).encode()
        )
        self.sensor_windows = defaultdict(lambda: deque(maxlen=window_size))
        self.sensor_stats = {}  # EWMA mean/std per sensor

    def process(self):
        for message in self.consumer:
            reading = message.value
            sensor_id = reading['sensor_id']
            value = reading['value']

            # Оновлюємо ковзне вікно
            self.sensor_windows[sensor_id].append(value)
            window = list(self.sensor_windows[sensor_id])

            # Детекція аномалій
            if len(window) >= 30:
                result = self.detect_anomaly(sensor_id, value, window)
                if result['anomaly']:
                    self.producer.send('iot-anomalies', result)

    def detect_anomaly(self, sensor_id, current_value, window):
        mean = np.mean(window)
        std = np.std(window)

        z_score = (current_value - mean) / (std + 1e-9)
        is_anomaly = abs(z_score) > 3.5

        return {
            'sensor_id': sensor_id,
            'value': current_value,
            'z_score': round(z_score, 2),
            'anomaly': bool(is_anomaly),
            'window_mean': round(mean, 3),
            'window_std': round(std, 3),
            'severity': 'critical' if abs(z_score) > 5 else 'warning'
        }

Чому контекстуальна аномалія важливіша?

Температура 85°C може бути нормальною о 14:00, але критичною о 3 годині ночі. Контекстуальні моделі враховують годину, день тижня та місяць для побудови baseline. Це знижує кількість хибних спрацювань у 3 рази порівняно зі звичайним z-score.

def contextual_anomaly_detection(sensor_id: str,
                                   current_value: float,
                                   timestamp: pd.Timestamp,
                                   historical_data: pd.DataFrame) -> dict:
    """
    Нормальний діапазон залежить від:
    - Час доби (година)
    - День тижня
    - Сезон (місяць)
    Baseline будується окремо для кожного контексту.
    """
    # Контекст поточного моменту
    context = {
        'hour': timestamp.hour,
        'day_of_week': timestamp.dayofweek,
        'month': timestamp.month
    }

    # Історичні дані у тому ж контексті
    context_data = historical_data[
        (historical_data['sensor_id'] == sensor_id) &
        (historical_data['hour'] == context['hour']) &
        (historical_data['day_of_week'] == context['day_of_week'])
    ]['value']

    if len(context_data) < 10:
        return {'status': 'insufficient_context_data'}

    context_mean = context_data.mean()
    context_std = context_data.std()
    context_z = (current_value - context_mean) / (context_std + 1e-9)

    return {
        'sensor_id': sensor_id,
        'value': current_value,
        'context': context,
        'context_mean': round(context_mean, 3),
        'context_z_score': round(context_z, 2),
        'contextual_anomaly': abs(context_z) > 3,
        'context_samples': len(context_data)
    }

Як відрізнити несправність датчика від аномалії процесу?

Якщо в групі з п'яти датчиків лише один показує аномалію — ймовірніше, зламався датчик. Якщо всі п'ять — проблема в процесі. Ми використовуємо правило одиничної аномалії: при частці аномальних датчиків менше 25% робимо висновок про несправність сенсора, при частці більше 60% — про фізичну аномалію.

def distinguish_sensor_vs_process_anomaly(sensor_group: dict,
                                           anomalous_sensor_id: str) -> dict:
    """
    Якщо лише один датчик з групи аномальний → швидше за все датчик зламаний.
    Якщо всі/більшість датчиків аномальні → процес аномальний.
    Застосовується коли декілька датчиків вимірюють одну фізичну зону.
    """
    anomaly_count = sum(1 for s_id, result in sensor_group.items()
                        if result.get('anomaly', False))
    total = len(sensor_group)

    group_anomaly_ratio = anomaly_count / total

    if group_anomaly_ratio <= 0.25:
        return {
            'conclusion': 'sensor_fault',
            'sensor_id': anomalous_sensor_id,
            'anomaly_ratio': group_anomaly_ratio,
            'action': 'replace_or_recalibrate_sensor',
            'process_alert': False
        }
    elif group_anomaly_ratio >= 0.6:
        return {
            'conclusion': 'process_anomaly',
            'anomaly_ratio': group_anomaly_ratio,
            'action': 'investigate_physical_process',
            'process_alert': True
        }
    else:
        return {
            'conclusion': 'uncertain',
            'anomaly_ratio': group_anomaly_ratio,
            'action': 'manual_investigation',
            'process_alert': True  # про всяк випадок
        }

Що таке Edge-інференс і навіщо він потрібен?

Для промислових мереж з обмеженим зв'язком затримка до хмари неприйнятна. Ми розгортаємо ONNX-моделі на ESP32 або Raspberry Pi. Інференс виконується локально, без мережевої затримки, зі споживанням менше 100 мс на одне передбачення.

import onnxruntime as ort
import numpy as np

class EdgeAnomalyDetector:
    """
    ONNX модель розгортається на edge пристрої.
    Інференс без хмари: критично для промислових мереж з обмеженим зв'язком.
    """
    def __init__(self, model_path: str, window_size: int = 30):
        self.session = ort.InferenceSession(model_path)
        self.window_size = window_size
        self.buffer = []
        self.threshold = 0.5

    def infer(self, new_value: float) -> dict:
        self.buffer.append(new_value)
        if len(self.buffer) > self.window_size:
            self.buffer.pop(0)

        if len(self.buffer) < self.window_size:
            return {'status': 'warming_up'}

        # Нормалізація
        arr = np.array(self.buffer, dtype=np.float32)
        arr = (arr - arr.mean()) / (arr.std() + 1e-9)
        input_tensor = arr.reshape(1, self.window_size, 1)

        result = self.session.run(None, {'input': input_tensor})[0]
        anomaly_score = float(result[0][0])

        return {
            'anomaly': anomaly_score > self.threshold,
            'score': round(anomaly_score, 3),
            'latency_ms': 'local'  # немає мережевої затримки
        }

Порівняння методів детекції

Метод Чутливість Хибні спрацювання Складність впровадження
Z-score (ковзне вікно) Середня Висока (викиди-одинаки) Низька — 1 день
Контекстуальний Z-score Висока Низька — враховує час/день Середня — 2-3 дні
ML-модель (LSTM/Transformer) Дуже висока Дуже низька Висока — 2-3 тижні

ML-детекція дає в 3 рази менше хибних спрацювань порівняно з пороговими методами. Ми обираємо підхід під ваші дані.

Порівняння протоколів IoT для передачі даних
Протокол Затримка Надійність Застосування
MQTT Низька Висока IoT-датчики
CoAP Низька Середня Обмежені пристрої
HTTP Висока Висока Моніторинг

Для більшості промислових сценаріїв ми рекомендуємо MQTT завдяки низькій затримці та вбудованій якості обслуговування (QoS).

Що входить до роботи

  1. Аудит — аналіз поточних джерел даних, частоти, якості та доступної інфраструктури.
  2. Проєктування — вибір стеку (Kafka / Flink / Spark / ONNX), схеми топиків та моделей.
  3. Розробка — написання пайплайну, навчання baseline та ML-моделі, налаштування алертів.
  4. Інтеграція — підключення MQTT-брокера, InfluxDB, Grafana.
  5. Розгортання — деплой на серверах або edge-пристроях (ESP32, Raspberry Pi).
  6. Документація — опис архітектури, інструкції оператора, керівництво з донавчання.
  7. Підтримка — навчання персоналу, гарантійний супровід 1 місяць.

Терміни та орієнтовна вартість

Терміни залежать від складності та обсягів: базовий пайплайн з z-score — від 2 тижнів, комплексне рішення з ML та edge — до 8 тижнів. Вартість розраховується індивідуально. Сертифіковані інженери з 5-річним досвідом у ML та IoT гарантують результат. Приклад економії: скорочення хибних спрацювань на 80% дозволяє заощадити до 2 млн грн на рік на обслуговуванні. Щоб отримати оцінку вашого проєкту, зв'яжіться з нами — ми проаналізуємо ваші дані та запропонуємо оптимальну архітектуру за 3 дні.

Інтеграція з платформами: AWS IoT Core, Azure IoT Hub, Yandex IoT Core, EdgeX Foundry. MQTT брокери: Eclipse Mosquitto, EMQ X. Зберігання та візуалізація: Grafana + InfluxDB.

Потрібна консультація? Опишіть ваше завдання — ми підберемо рішення за один день.

Виявлення аномалій: автоенкодери, Isolation Forest, PyOD

Ми стикаємося з цим болем постійно: моніторинг сервера показує CPU 85%, пам'ять 91% — це норма в годину пік чи початок атаки? Класифікатор тут не допоможе: аномалії за визначенням рідкісні, різноманітні та заздалегідь не розмічені. Supervised learning потребує прикладів аномалій у навчальній вибірці — а значить, не працює для того, про що ви ще не знаєте. Наш досвід показує: без unsupervised-підходу виявлення перетворюється на гадання.

Чому виявлення аномалій потребує unsupervised підходу?

Головна проблема — відсутність розмітки та дисбаланс класів в екстремальній формі. Фрод-транзакції становлять 0.01–0.1% від загального об'єму. Виробничий дефект — 0.5–3%. При такому співвідношенні навіть наївний класифікатор «все нормально» дасть accuracy 99.9% і precision/recall для аномального класу, близькі до нуля. Supervised-моделі тут безсилі.

Друга проблема — «нормальність» завжди контекстна. Чи нормально, що користувач логіниться о 3 годині ночі? Залежить від його історії та часової зони. Чи нормальна вібрація підшипника 2.3 мм/с? Залежить від режиму роботи верстата та його віку. Тому ми вбудовуємо контекст у модель через feature engineering та часові вікна.

Третя — оцінка якості. Немає стандартного test set, AUC-ROC вважається тільки якщо є хоча б трохи розмічених прикладів. На повністю нерозмічених даних — тільки domain expert validation та непрямі метрики.

Як відрізнити аномалію від шуму в реальному часі?

Відповідь — адаптивні пороги та моніторинг статистик моделі. У розділі кейсу покажемо, як це працює.

Методи та інструменти

Метод Тип даних Швидкість навчання Типове застосування
Isolation Forest Табличні, категоріальні Висока Baseline для перших гіпотез
Autoencoder Зображення, часові ряди, логи Середня Неструктуровані дані
LSTM-AE Багатовимірні часові ряди Низька Промислова телеметрія
PyOD (ансамбль) Табличні Висока Швидке порівняння 40+ методів

Isolation Forest — стандартний baseline для табличних даних. Ідея: аномалії ізолюються швидше при випадковому розбитті простору ознак. Працює добре при contamination 0.01–0.1, стійкий до масштабу ознак, не потребує нормалізації. Реалізація в sklearn.ensemble.IsolationForest.

Типова помилка: ставити contamination='auto' без розуміння даних. Auto-режим передбачає поріг -0.5, що не завжди відповідає реальній частці аномалій. Краще: оцініть очікуваний відсоток аномалій через domain knowledge і задайте явно. Ми гарантуємо підбір contamination під ваш кейс.

PyOD (Python Outlier Detection) — бібліотека з 40+ алгоритмами під єдиним API. Включає: OCSVM, LOF, COPOD, ECOD, DeepSVDD, AutoEncoder. Зручно для швидкого порівняння методів на одних даних.

Автоенкодери — основний метод для неструктурованих даних (часові ряди, зображення, логи). Ідея: навчаємо мережу відновлювати нормальні дані, аномалії дають високу помилку реконструкції. Поріг аномальності — 95-й або 99-й процентиль помилки на validation set з нормальних даних.

Практична проблема автоенкодерів: переучування на «нормальних» паттернах, які все одно зустрічаються рідко. Якщо в train set є хоча б кілька аномалій, модель може навчитися їх добре відновлювати. Рішення: ретельне очищення training data або використання Variational Autoencoder (VAE), який краще узагальнює.

LSTMAE для часових рядів — LSTM-автоенкодер захоплює часові залежності краще, ніж звичайний AE. Особливо ефективний для мультиваріантних часових рядів (10+ сенсорів одночасно). Реалізація через PyTorch, навчання з MSELoss на ковзних вікнах.

Детально: виявлення аномалій у промислових часових рядах

Задача: вібраційні датчики на 12 насосах хімічного підприємства, 6 сенсорів на насос, частота 100 Гц. Потрібно попередити про наближену поломку за 4–24 години.

Архітектура рішення:

Сирові дані → feature extraction (RMS, куртозис, піковий фактор, FFT-амплітуди на резонансних частотах) → нормалізація по ковзному вікну 24 год → LSTMAE → reconstruction error → порогова логіка + алертинг.

Розмір вікна LSTM: 60 секунд (6000 точок на 100 Гц). Занадто мале вікно — не захоплює повільні паттерни. Занадто велике — втрачає чутливість до швидких змін.

Поріг аномальності: не фіксований, а адаптивний. threshold = mean(errors_last_7d) + 3 * std(errors_last_7d). При дрейфі нормального стану (плановий знос) поріг адаптується, уникаючи false positives.

Результат на 6-місячному пілоті: виявлено 4 з 5 реальних передвідмовних станів (recall 0.8), 2 хибні тривоги за 6 місяців (precision 0.67). До впровадження: 3 незаплановані зупинки зі значними збитками. Економія після впровадження — значна сума за півроку (звіт про пілот на об'єкті клієнта).

Фрод-детекція: специфіка фінансових даних

Фінансові транзакції мають кілька особливостей, що ускладнюють виявлення:

  • Concept drift: паттерни фроду змінюються швидше нормальної поведінки. Модель, навчена півроку тому, застаріває.
  • Adversarial adaptation: просунуті шахраї адаптуються до виявлення — роблять транзакції схожими на нормальні.
  • Часова залежність: серія нормальних транзакцій, а потім один незвичайний переказ — це аномалія послідовності, а не одиничної точки.

Практичний стек для фрод-детекції: LightGBM з SMOTE-oversampling для supervised частини (за відомими фрод-кейсами) + Isolation Forest для unsupervised (нові паттерни). Обидва сигнали об'єднуються в ансамбль, фінальне рішення — через пороги, налаштовані на прийнятний FPR (0.1–1% від транзакцій на ручну перевірку).

Як оцінити якість без розмітки?

Коли ground truth немає, для оцінки використовуємо:

  • Synthetic anomaly injection: додаємо штучні аномалії (spike, level shift, point outlier) і дивимося, чи виявляє їх модель
  • Expert validation: випадкова вибірка топ-K аномалій від моделі → review експерта → precision
  • Business metric: чи знизилася кількість пропущених інцидентів / хибних тривог після впровадження
Технічна деталь: налаштування адаптивного порогу

Поріг обчислюється як mean(errors) + k * std(errors) на ковзному вікні 7 днів. Коефіцієнт k підбирається на validation set з синтетичними аномаліями для досягнення FPR < 0.1%. При дрейфі ознак вікно автоматично зсувається.

Процес роботи

  1. Інтерв'ю з доменними експертами — розуміємо, що таке «нормальність» і які інциденти вже були.
  2. EDA та підготовка даних — очищення, створення ознак, часові вікна.
  3. Baseline (Isolation Forest) — швидка валідація на відомих інцидентах.
  4. Вибір та кастомізація моделі — Autoencoder / LSTM-AE / ансамбль.
  5. Навчання, валідація з синтетичними аномаліями.
  6. Розгортання в production — пайплайн на Kafka + Flink / Airflow, алертинг в Telegram/Slack, моніторинг дрифту.
  7. Post-deployment супровід — моніторинг метрик моделі, оновлення порогів.

Що входить у роботу

  • Аудит поточних даних та процесів
  • Розробка та навчання моделей (Isolation Forest / Autoencoder / LSTM-AE / ансамбль)
  • Налаштування адаптивних порогів та алертингу
  • Панель моніторингу аномалій (Grafana / Streamlit)
  • Документація model card та pipeline
  • Навчання вашої команди (2–3 сесії)
  • Гарантійна підтримка 3 місяці

Терміни: baseline-система з одним методом — 2–4 тижні. Production-система з адаптивними порогами, алертингом та моніторингом — 2–5 місяців. Вартість розраховується індивідуально під ваш кейс.

Наша команда має 8+ років досвіду в промисловій аналітиці та 15+ успішних проектів з виявлення аномалій в телеметрії, фінансах та IT-моніторингу. Отримайте консультацію — розкажемо, як вирішити вашу задачу.