Расскажи про опыт работы с Apache Kafka в MLOps

«Расскажи про опыт работы с Apache Kafka в MLOps» — вопрос из категории MLOps и деплой моделей, который задают на 26% собеседований Data Scientist / ML Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

В контексте MLOps я использовал Apache Kafka как бэкбон для потоковой обработки данных и событийного управления жизненным циклом моделей. Основные сценарии: доставка функций (features) для онлайн-инференса и потоковый мониторинг дрейфа данных.

Архитектурный пример — потоковый инференс:

  1. Продюсеры (микросервисы или IoT-устройства) публикуют сырые данные или события в топик Kafka, например, raw-transactions.
  2. Stream-процессор (написанный с помощью confluent-kafka и scikit-learn/PyFunc MLflow) подписывается на этот топик, применяет пайплайн преобразования признаков и загруженную ML-модель для предсказания.
  3. Результаты записываются в топик predictions, откуда их потребляют другие сервисы.

Пример потребителя на Python для мониторинга:

from confluent_kafka import Consumer, KafkaError
import pandas as pd
from evidently.report import Report
from evidently.metrics import DataDriftTable

conf = {'bootstrap.servers': 'kafka-broker:9092',
        'group.id': 'ml-monitoring-group',
        'auto.offset.reset': 'latest'}
consumer = Consumer(conf)
consumer.subscribe(['model-input-features'])

batch = []
while True:
    msg = consumer.poll(1.0)
    if msg is None:
        continue
    if msg.error():
        print(f"Consumer error: {msg.error()}")
        continue

    feature_record = json.loads(msg.value().decode('utf-8'))
    batch.append(feature_record)

    if len(batch) >= 1000:  # Анализируем батч
        current_df = pd.DataFrame(batch)
        # Сравниваем с референсным датасетом (например, обучающим)
        drift_report = Report(metrics=[DataDriftTable()])
        drift_report.run(reference_data=ref_df, current_data=current_df)
        if drift_report.show()['metrics'][0]['result']['dataset_drift']:
            # Триггер на переобучение или оповещение
            trigger_retraining_alert()
        batch = []

Ключевые настройки в MLOps:

  • Сериализация: Использование Avro или Protobuf через Schema Registry для строгой контрактности данных между сервисами.
  • Ретеншен топиков: Настройка политик хранения для топиков с сырыми данными и предсказаниями, чтобы можно было повторно проиграть события для отладки или переобучения модели.
  • Интеграция с пайплайнами: Запуск переобучения модели как реакции на событие в Kafka (например, при обнаружении дрейфа).

Таким образом, Kafka выступает центральной нервной системой для асинхронной, отказоустойчивой и масштабируемой ML-инфраструктуры.