Расскажи про опыт работы с Clickstream

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

Ответ

Я работал с обработкой Clickstream-данных для построения аналитических дашбордов и систем рекомендаций. Clickstream — это последовательность событий (просмотры, клики, скроллы) пользователя на сайте или в приложении.

Мой стек и задачи:

  1. Сбор и инжекция: Настраивал отправку событий с фронтенда (через JavaScript SDK или Pixel) в Kafka. Использовал схему Avro для сериализации данных, чтобы обеспечить консистентность и эффективность.

    // Пример события в Kafka
    {
     "user_id": "u123",
     "session_id": "sess_abc",
     "event_type": "page_view",
     "page_url": "/product/123",
     "timestamp": 1678886400000
    }
  2. Обработка (Streaming & Batch):

    • Реал-тайм пайплайны (Kafka Streams/Flink): Для агрегации ключевых метрик (например, активные пользователи за последние 5 минут) и обнаружения аномалий.
    • Батч-обработка (Spark): Для глубокого анализа: построение воронок, вычисление долгосрочных метрик (LTV, retention).

      # PySpark: Расчет средней глубины просмотра за сессию
      from pyspark.sql import Window
      from pyspark.sql import functions as F
      
      window_spec = Window.partitionBy("session_id").orderBy("timestamp")
      df_with_seq = df.withColumn("event_seq", F.row_number().over(window_spec))
      avg_session_depth = df_with_seq.groupBy("session_id").agg(F.max("event_seq").alias("depth"))
  3. Хранение и анализ: Загружал обработанные данные в ClickHouse для быстрой аналитики (ad-hoc запросы) и в S3 (в формате Parquet) как data lake для исторических данных и ML-моделей.

  4. Сложности и решения:

    • Дубликаты и потерянные события: Решал через идемпотентность обработки в Kafka и дедупликацию по event_id на этапе загрузки в хранилище.
    • Семантика сессий: Реализовывал логику определения сессий (по таймауту) как в потоковой обработке, так и в SQL (используя LAG и условную агрегацию).
    • Volume: Оптимизировал партиционирование в S3 по дате (dt=2023-01-01/) и использовал сжатие (Snappy, Zstd).

Результатом работы были дашборды в Tableau/Grafana для продукт-менеджеров и автоматические A/B-тесты.