Какие инструменты используются для очистки событий?

«Какие инструменты используются для очистки событий?» — вопрос из категории Качество данных, который задают на 33% собеседований Data Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

Под "очисткой событий" (event cleansing) обычно понимается валидация и коррекция данных в потоковых или пакетных событиях (например, клики, просмотры). В моих проектах для этого использовались следующие подходы и инструменты:

  1. Валидация схемы: Проверка, что событие содержит все обязательные поля с правильными типами. Использовал Apache Avro с заранее определённой схемой (Schema Registry) или библиотеку Pydantic в Python для валидации на лету.

    from pydantic import BaseModel, ValidationError
    from datetime import datetime
    
    class ClickEvent(BaseModel):
        user_id: int
        url: str
        timestamp: datetime
        ip_address: str
    
    try:
        valid_event = ClickEvent(**raw_event_data)
    except ValidationError as e:
        # Отправить событие в DLQ (Dead Letter Queue) для дальнейшего разбора
        send_to_dlq(raw_event_data)
  2. Очистка и нормализация:

    • Приведение значений к единому формату (например, приведение всех дат к UTC, нормализация строк). Делал это с помощью функций в PySpark или Pandas.
    • Исправление опечаток в категориальных полях (например, названия городов) с помощью алгоритмов нечёткого сравнения, таких как fuzzywuzzy.
  3. Фильтрация аномалий:

    • Статистические методы (например, отсечение событий по межквартильному размаху - IQR) для числовых полей (длительность сессии, сумма заказа).
    • Правила на основе доменного знания (например, отфильтровать события с будущими временными метками).
  4. Инструменты: Основная логика очистки часто писалась на Python (Pandas/PySpark). Для оркестрации пайплайнов очистки использовался Apache Airflow. Проблемные события отправлялись в Dead Letter Queue (например, в Kafka-топик или S3-бакет) для ручного аудита.