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

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

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

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

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

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

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 дні): розгортаємо в продакшен, налаштовуємо дашборди та алерти.

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

Чому дата-інжиніринг визначає успіх ML-моделі

Минулого року до нас звернулася компанія, яка витратила $50 000 на навчання NLP-моделі, але отримала лише 60% точності на продакшені. Причина — data leakage через випадковий split часових даних. Перед тим як навчати модель, потрібно зрозуміти структуру даних: чи є дублі, як часто змінюється схема, наскільки репрезентативна вибірка. Дата-інжиніринг для ML — це не просто ETL, а побудова відтворюваної інфраструктури, яка робить навчання надійним, а перенавчання — передбачуваним. За досвідом нашої команди (понад 8 років у дата-інжинірингу, 30+ проектів у ML) кожна друга проблема в продакшені пов’язана не з архітектурою моделі, а з якістю даних. Замовте аудит ваших даних — оцінимо поточний пайплайн безкоштовно.

Як ETL-пайплайни для ML відрізняються від BI

ETL для аналітики та ETL для ML — різні завдання. В аналітиці важлива агрегація, у ML — індивідуальні записи з історією. В аналітиці train/val/test split не потрібен, у ML — критичний. В аналітиці skew даних заважає інтерпретації, у ML — безпосередньо впливає на якість моделі.

Інструменти. Apache Spark для великих обсягів (10GB+): PySpark з DataFrames, оптимізації через partitioning та caching. dbt для трансформацій поверх DWH (Snowflake, BigQuery, Redshift) — декларативно, версіонується, тестується. Pandas + Polars для обсягів до кількох GB — Polars у 5–10x швидше за Pandas на типових трансформаціях.

Temporal splits. Для ML важливо, що split за часом, а не випадковий. Якщо дані часові (транзакції, події користувачів), випадковий split дає data leakage: модель бачить «майбутні» дані при навчанні. Правило: train на періоді T1–T2, validation на T2–T3 (з gap для запобігання leakage), test на T3–T4. Неправильний split може коштувати 10–15% якості моделі на валідації. Temporal split best practices (scikit-learn docs)

Інкрементальні пайплайни. Модель перенавчається щотижня на нових даних. Потрібен пайплайн, який інкрементально додає нові записи до навчальної вибірки, не перевантажуючи все з нуля. Delta Lake або Apache Iceberg — формати з ACID-транзакціями, Change Data Capture, time travel.

Як уникнути training-serving skew за допомогою Feature Store

Feature Store вирішує проблему розсинхронізації між навчанням та інференсом. Найпідступніша помилка в ML-інфраструктурі — training-serving skew: ознака обчислюється по-різному в навчанні та в продакшені. Модель вчиться на «правильних» даних, а інференс отримує інші.

Feast (open source) — офлайн store на Parquet/Delta в S3 для навчання, онлайн store на Redis для low-latency інференсу (<10ms). Feature definitions як Python-код:

from feast import FeatureView, Field
from feast.types import Float32, Int64

user_features = FeatureView(
    name="user_features",
    entities=["user_id"],
    schema=[
        Field(name="purchase_count_7d", dtype=Int64),
        Field(name="avg_session_duration", dtype=Float32),
    ],
    ttl=timedelta(days=7),
    source=user_features_source,
)

Один definition використовується всюди — немає розбіжностей.

Потокові ознаки. Коли ознака має оновлюватися в реальному часі (кількість транзакцій за останні 10 хвилин), потрібна потокова обробка. Apache Kafka + Apache Flink або Kafka Streams для обчислення ознак у реальному часі → запис в онлайн store. Складніше, дорожче, потрібно лише коли staleness ознак критична для якості.

Розмітка даних: як не витратити бюджет даремно

Розмітка — найтрудомісткіша та недооцінювана частина ML-проекту. Погано розмічені дані не виправить жодна архітектура.

Label Studio — open source, підтримує розмітку зображень (bounding box, polygon, segmentation), тексту (NER, класифікація), аудіо, відео. Піднімається за 10 хвилин через Docker. Для невеликих команд — перший вибір.

Оцінка якості розмітки. Inter-annotator agreement — наскільки згодні розмітники між собою. Cohen's Kappa > 0.8 — добре, 0.6–0.8 — прийнятно, < 0.6 — завдання неоднозначне або інструкція погана. Перетин розміток (10–20% прикладів розмічають два незалежних анотатори) — обов'язкова практика.

Active learning. Не розмічати випадкові приклади, а вибирати ті, на яких модель найбільш невпевнена (low confidence, high uncertainty). Дозволяє досягти тієї ж якості при 50–70% обсягу розмітки. Modals, Prodigy, Label Studio підтримують active learning workflows. На одному з проектів для NLP ми скоротили бюджет на розмітку в 2,5 рази завдяки active learning — економія склала $15 000 на 100 000 розмічених прикладів.

Синтетичні дані. Коли реальних даних мало або отримати їх дорого. Для CV: рендеринг у Blender/Unity з реалістичними текстурами (domain randomization). Для NLP: parafrase через LLM, backtranslation. Ризик: модель навчається на distribution синтетичних даних, а не реальних — потрібна обережність і перевірка на реальному holdout.

Якість даних: валідація та моніторинг

Great Expectations — de facto стандарт для data validation у ML-пайплайнах. Expectations — це декларативні твердження про дані: «колонка age містить значення від 0 до 120», «колонка user_id не містить null», «розподіл amount не відхиляється більш ніж на 20% від baseline». Запускається в пайплайні, при провалі — блокує проходження.

Pandera — Pythonic alternative для pandas/polars DataFrames. Schema-based validation з type hints:

import pandera as pa

schema = pa.DataFrameSchema({
    "user_id": pa.Column(int, nullable=False),
    "score": pa.Column(float, pa.Check.between(0, 1)),
    "label": pa.Column(str, pa.Check.isin(["positive", "negative", "neutral"])),
})

Data freshness. Модель очікує дані за останні N днів. ETL впав, дані не оновилися — модель використовує застарілі ознаки. Моніторинг свіжості даних: timestamp останнього запису в кожній таблиці, алерт при затримці > порога.

Дедуплікація. Дублікати в навчальній вибірці завищують метрики (одні й ті самі приклади в train і val) і спотворюють ваги моделі. MinHash LSH для наближеної дедуплікації великих датасетів. Для точної — хеш за нормалізованим контентом.

Інструмент Область застосування Коли вибирати
Great Expectations Універсальна, таблиці, пайплайни Великі команди, багато метаданих
Pandera pandas/polars DataFrames Python-centric проекти, type hints
Deequ Apache Spark, великі дані Якщо пайплайн вже на Spark

Сховища та формати

Формат Найкраще для Особливості
Parquet Батчеве навчання, аналітика Columnar, ефективне стиснення
Delta Lake Інкрементальні апдейти, ACID Time travel, schema evolution
Apache Iceberg Enterprise, multi-engine Найкращий catalog, hidden partitioning
HDF5 Числові масиви (CV датасети) Ієрархічна структура
TFDS / datasets Стандартизовані ML датасети Hugging Face datasets — зручний для NLP

Для більшості ML-проектів на старті: Parquet в S3 + DVC для версіонування. Delta Lake або Iceberg — коли з'являється потреба в інкрементальних оновленнях або time travel.

Типові помилки при побудові пайплайнів

  • Пропуск перевірки свіжості даних. Якщо ETL падає вночі, а модель запускається вранці — вона отримує дані 24-годинної давності. Рішення: алерт при затримці > 30 хвилин.
  • Відсутність версіонування даних. Не можна відтворити експеримент, бо дані змінилися. DVC або Delta Lake time travel виправляють це.
  • Забувають про schema evolution. Нове поле з’являється, а пайплайн падає. Автоматичне виявлення змін схеми через Great Expectations.

Active learning дозволяє скоротити бюджет на розмітку до 50–70%. На одному проекті це склало економію $15 000 на 100 000 розмічених прикладів. Закажіть консультацію — розрахуємо потенційну економію для вашого кейсу.

Що входить у проект з дата-інжинірингу для ML

Ми надаємо повний цикл:

  • Аудит існуючих даних та пайплайнів (1 тиждень).
  • Проектування архітектури: вибір інструментів, форматів, способів розмітки.
  • Реалізація ETL/ELT пайплайну з валідацією та моніторингом.
  • Документація коду та процесів (model card, data card).
  • Навчання вашої команди роботі з пайплайном.
  • SLA на супровід та підтримку.

Терміни: від 2 до 6 тижнів залежно від обсягу даних і складності інтеграцій.

Як ми будуємо пайплайн: покроково

  1. Аудит існуючих даних. Профілювання: ydata-profiling (колишній pandas-profiling) генерує HTML-репорт зі статистиками, дистрибуціями, кореляціями, missing values за хвилини.
  2. Проектування пайплайну. Визначаємо джерела даних, частоту оновлення, вимоги до latency ознак, обсяги.
  3. Реалізація та тестування. Unit-тести на трансформації, integration-тести на пайплайн, data validation через Great Expectations.
  4. Деплой та моніторинг. Алерти на freshness, quality checks, аномалії в обсягах даних.

Чому варто довірити це нам

Ми займаємося дата-інжинірингом та ML з понад 8-річним досвідом. За цей час реалізували понад 40 проектів — від побудови пайплайнів для NLP-моделей до розмітки датасетів для комп’ютерного зору. Гарантуємо відтворюваність пайплайнів та повну прозорість процесів. У кожному проекті використовуємо інструменти з відкритим кодом, щоб ви не були прив’язані до вендора.

Зв’яжіться з нами для безкоштовного аудиту ваших даних — оцінимо поточний пайплайн і запропонуємо roadmap. Замовте побудову ML-пайплайну під ключ.