На производственной линии датчик температуры резко показывает 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).
Что входит в работу
- Аудит — анализ текущих источников данных, частоты, качества и доступной инфраструктуры.
- Проектирование — выбор стека (Kafka / Flink / Spark / ONNX), схемы топиков и моделей.
- Разработка — написание пайплайна, обучение baseline и ML-модели, настройка алертов.
- Интеграция — подключение MQTT-брокера, InfluxDB, Grafana.
- Развёртывание — деплой на серверах или edge-устройствах (ESP32, Raspberry Pi).
- Документация — описание архитектуры, инструкции оператора, руководство по дообучению.
- Поддержка — обучение персонала, гарантийное сопровождение 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.
Нужна консультация? Опишите вашу задачу — мы подберём решение за один день.







