Потокові ML-пайплайни для real-time інференсу на Kafka, Flink та ONNX

Fraud detection в реальному часі — типовий кейс, де latency > 100 мс означає втрату грошей. Ми розробляємо потокові ML-пайплайни, які обробляють події із затримкою до 50 мс, використовуючи Kafka, Flink та ONNX Runtime. Один із проектів — система антифроду для фінтех-компанії, що обробляє 50 000 тран

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

Часті запитання

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

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

Fraud detection в реальному часі — типовий кейс, де latency > 100 мс означає втрату грошей. Ми розробляємо потокові ML-пайплайни, які обробляють події із затримкою до 50 мс, використовуючи Kafka, Flink та ONNX Runtime. Один із проектів — система антифроду для фінтех-компанії, що обробляє 50 000 транзакцій на секунду з p99 latency 45 мс. Інший приклад — пайплайн для платіжного агрегатора з throughput 20K подій/с та латентністю p99 55 мс, що скоротило chargeback на 35%. Перехід з batch-пайплайну на потоковий дозволив цьому клієнту знизити витрати на інфраструктуру на 40% — окупність проекту склала 3 місяці.

Ми реалізували понад 50 подібних рішень. Apache Flink — один із ключових інструментів, що забезпечує exactly-once семантику та відмовостійкість. Сертифіковані інженери гарантують стабільність та масштабованість пайплайну.

Проблеми, які вирішуємо

Латентність. Класичні batch-пайплайни вносять затримки від хвилин до годин. Для real-time скорингу транзакцій або оновлення рекомендацій це неприйнятно. Ми використовуємо sliding window агрегації та online feature store (Redis) з lookup < 5ms, щоб ознаки були доступні миттєво.

State management. Потокова обробка вимагає консистентного стану при збоях та перезапусках. Apache Flink надає exactly-once семантику та checkpointing, що критично для фінансових застосувань.

Model versioning. При оновленні моделі не можна зупиняти пайплайн. Реалізуємо A/B тестування через тіньовий трафік та поступову розкатку нових версій за допомогою feature flags.

Архітектура потокового ML-пайплайну

[Kafka / Kinesis / Pulsar] ↓ [Feature Computation] ← Flink / Spark Streaming / Kafka Streams (агрегації, вікна, joins) ↓ [Feature Store Online] ← Redis / DynamoDB (< 5ms lookup) ↓ [Model Inference] ← Triton / TorchServe / ONNX Runtime (< 20ms) ↓ [Decision Engine] ← бізнес-правила + ML score ↓ [Action / Output Kafka] ← downstream системи 

Чому Kafka Streams і ONNX?

Kafka Streams вбудовується в будь-який мікросервіс і не потребує окремого кластера для обробки. Для високонавантажених сценаріїв використовуємо Flink з паралелізмом до 16. ONNX Runtime дозволяє виконувати моделі на CPU з latency < 5ms та підтримує квантизацію INT8 для додаткового прискорення.

Як забезпечити exactly-once семантику?

Apache Flink підтримує exactly-once за рахунок checkpointing та узгоджених snapshot'ів стану. В документації Apache Flink описані механізми, які ми застосовуємо в продакшені. Це гарантує, що при збоях жодна подія не буде втрачена або оброблена двічі.

Реалізація потокового пайплайну

Обчислення ознак в реальному часі

from confluent_kafka import Consumer, Producer import json import redis import numpy as np import time from collections import deque, defaultdict import threading class StreamFeatureComputer: """Обчислення ознак в реальному часі""" def __init__(self, kafka_config: dict, redis_url: str): self.consumer = Consumer(kafka_config) self.producer = Producer({'bootstrap.servers': kafka_config['bootstrap.servers']}) self.redis = redis.from_url(redis_url) self.window_store = defaultdict(lambda: deque(maxlen=1000)) def compute_user_features(self, user_id: str, event: dict) -> dict: """Online-ознаки для користувача""" key_prefix = f"user:{user_id}" now = event['timestamp'] # Sliding window агрегації через Redis pipe = self.redis.pipeline() # Транзакційні ознаки event_key = f"{key_prefix}:events" pipe.lpush(event_key, json.dumps({ 'amount': event.get('amount', 0), 'ts': now, 'type': event.get('type', 'unknown') })) pipe.ltrim(event_key, 0, 999) # Тримаємо останні 1000 подій pipe.expire(event_key, 86400) # TTL 24 години pipe.execute() # Агрегації за різні вікна raw_events = self.redis.lrange(event_key, 0, -1) events = [json.loads(e) for e in raw_events] # Сортування за часом events.sort(key=lambda x: x['ts'], reverse=True) window_1h = [e for e in events if now - e['ts'] <= 3600] window_24h = [e for e in events if now - e['ts'] <= 86400] amounts_1h = [e['amount'] for e in window_1h] amounts_24h = [e['amount'] for e in window_24h] features = { 'user_id': user_id, 'tx_count_1h': len(window_1h), 'tx_count_24h': len(window_24h), 'tx_amount_sum_1h': sum(amounts_1h), 'tx_amount_sum_24h': sum(amounts_24h), 'tx_amount_avg_1h': np.mean(amounts_1h) if amounts_1h else 0, 'tx_amount_max_1h': max(amounts_1h) if amounts_1h else 0, 'tx_amount_std_1h': np.std(amounts_1h) if len(amounts_1h) > 1 else 0, 'unique_merchants_1h': len(set(e.get('merchant_id') for e in window_1h)), 'time_since_last_tx': now - events[0]['ts'] if events else 9999, } return features def compute_velocity_features(self, entity_id: str, event_type: str, windows: list[int] = [60, 300, 3600]) -> dict: """Velocity checks: частота подій за різні вікна""" features = {} now = int(time.time()) for window in windows: key = f"velocity:{entity_id}:{event_type}:{window}" # Increment та expire pipe = self.redis.pipeline() pipe.incr(key) pipe.expire(key, window) count, _ = pipe.execute() features[f"count_{window}s"] = count return features 

Потоковий інференс з ONNX

import onnxruntime as ort import asyncio from aiohttp import ClientSession class StreamMLInference: """Низьколатентний інференс у потоці""" def __init__(self, model_path: str, feature_store: redis.Redis): # ONNX для максимальної швидкості opts = ort.SessionOptions() opts.inter_op_num_threads = 2 opts.intra_op_num_threads = 2 opts.graph_optimization_level = ort.GraphOptimizationLevel.ORT_ENABLE_ALL self.session = ort.InferenceSession( model_path, sess_options=opts, providers=['CPUExecutionProvider'] ) self.feature_store = feature_store self.input_names = [inp.name for inp in self.session.get_inputs()] def predict(self, features: dict) -> dict: """Інференс < 5ms для tabular моделі""" # Формування input tensor feature_vector = np.array([[features.get(name, 0.0) for name in self.input_names]], dtype=np.float32) start = time.perf_counter() outputs = self.session.run(None, {self.input_names[0]: feature_vector}) latency_ms = (time.perf_counter() - start) * 1000 score = float(outputs[0][0][1]) # Probability of positive class return { 'score': score, 'decision': 'block' if score > 0.8 else 'review' if score > 0.5 else 'allow', 'latency_ms': latency_ms } def batch_predict(self, features_list: list[dict]) -> list[dict]: """Батч-інференс для мікробатчів""" if not features_list: return [] feature_matrix = np.array([[f.get(name, 0.0) for name in self.input_names] for f in features_list], dtype=np.float32) outputs = self.session.run(None, {self.input_names[0]: feature_matrix}) scores = outputs[0][:, 1].tolist() return [ {'score': s, 'decision': 'block' if s > 0.8 else 'review' if s > 0.5 else 'allow'} for s in scores ] 

Apache Flink пайплайн (Python API)

from pyflink.datastream import StreamExecutionEnvironment from pyflink.datastream.connectors.kafka import KafkaSource, KafkaSink from pyflink.common import WatermarkStrategy, Types from pyflink.datastream.window import TumblingEventTimeWindows, SlidingEventTimeWindows from pyflink.common.time import Time def build_flink_ml_pipeline(): env = StreamExecutionEnvironment.get_execution_environment() env.set_parallelism(4) # Kafka source source = KafkaSource.builder() \ .set_bootstrap_servers("kafka:9092") \ .set_topics("transactions") \ .set_group_id("ml-pipeline") \ .set_value_only_deserializer(JsonRowDeserializationSchema()) \ .build() stream = env.from_source( source, WatermarkStrategy.for_monotonous_timestamps(), "Kafka Source" ) # Обчислення агрегатів за 5-хвилинне ковзне вікно windowed = stream \ .key_by(lambda event: event['user_id']) \ .window(SlidingEventTimeWindows.of(Time.minutes(5), Time.seconds(30))) \ .aggregate(TransactionAggregator()) # Приєднання до static features з бази enriched = windowed.map(EnrichWithStaticFeatures()) # ML інференс scored = enriched.map(MLScoringFunction()) # Sink: дії в реальному часі sink = KafkaSink.builder() \ .set_bootstrap_servers("kafka:9092") \ .set_record_serializer(JsonRowSerializationSchema("ml-decisions")) \ .build() scored.sink_to(sink) env.execute("ML Streaming Pipeline") 

Моніторинг та метрики

class StreamPipelineMonitor: """Метрики для real-time ML пайплайну""" def __init__(self, prometheus_port: int = 8000): from prometheus_client import Counter, Histogram, Gauge, start_http_server self.events_processed = Counter('ml_events_total', 'Total events processed', ['decision']) self.inference_latency = Histogram('ml_inference_latency_ms', 'Inference latency in milliseconds', buckets=[1, 5, 10, 20, 50, 100, 500]) self.feature_lag = Gauge('feature_store_lag_ms', 'Time between event and feature availability') self.model_score_dist = Histogram('ml_model_score', 'Distribution of model scores', buckets=[0.1*i for i in range(11)]) start_http_server(prometheus_port) def record_inference(self, result: dict): self.events_processed.labels(decision=result['decision']).inc() self.inference_latency.observe(result.get('latency_ms', 0)) self.model_score_dist.observe(result['score']) 

Порівняння підходів до потокового інференсу

Підхід Латентність (p99) Масштабування Складність впровадження
Kafka Streams + ONNX < 50ms Горизонтальне, до 100K ev/s Середня
Apache Flink + Triton < 80ms До 500K ev/s, stateful Висока
Spark Streaming + TensorFlow < 200ms До 1M ev/s, мікробатчі Середня

Порівняння online feature store

Рішення Latency lookup Масштабування Ціна
Redis < 1ms До 100K ops/s Низька
DynamoDB < 5ms Автоматичне Середня
Aerospike < 1ms До 1M ops/s Висока
Типові помилки при побудові потокового ML
  • Відсутність watermarking: події із затримкою можуть спотворити агрегації. Завжди налаштовуйте allowed lateness.
  • Ігнорування backpressure: використовуйте reactive streams або динамічний parallelism.
  • Збереження стану лише в пам'яті: обов'язково використовуйте checkpointing та реплікацію state backend.
  • Прямий виклик моделей у потоці: краще винести інференс в окремий мікросервіс з чергою.
  • Відсутність моніторингу: Prometheus + Grafana для latency, throughput, error rate.

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

Базова реалізація пайплайну (Kafka + Flink + Feature Store + ONNX) займає 3-4 тижні. Складні сценарії з кастомними агрегаціями та A/B тестуванням — до 6 тижнів. Ми підбираємо оптимальний стек під ваш сценарій. Зв'яжіться — оцінимо проект протягом 2 днів.

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

  • Архітектурна схема
  • Код пайплайну з коментарями
  • Дашборди моніторингу (Prometheus/Grafana)
  • Документація
  • Навчання команди
  • Підтримка 1 місяць після запуску

Результати та економіка

Економічна ефективність

Перехід з batch-пайплайну на потоковий дозволяє скоротити витрати на інфраструктуру на 30-40% за рахунок відмови від проміжного зберігання даних. Середня окупність проекту — 3 місяці. Наприклад, один із клієнтів після впровадження знизив chargeback на 35% та скоротив час детекції шахрайства з 5 хвилин до 50 мс.

Покроковий план впровадження

  1. Аналітика та аудит (2-3 дні): вивчаємо поточні дані, вимоги до latency та throughput, вибираємо стек.
  2. Проектування архітектури (3-5 днів): розробляємо схему пайплайну, визначаємо контракти.
  3. Реалізація (1-2 тижні): пишемо код потокової обробки, feature engineering, інтеграцію з ML-моделлю.
  4. Тестування (3-5 днів): навантажувальне тестування, перевірка відмовостійкості, оптимізація.
  5. Деплой та моніторинг (2-3 дні): розгортаємо в продакшен, налаштовуємо дашборди та алерти.

Замовте аудит поточної системи та отримайте комерційну пропозицію. Зв'яжіться з нами для обговорення вашого кейсу.