Сталкивался ли с объединением потоков данных в реальном времени (stream-stream joins)?

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

Ответ

Да, реализовывал stream-stream joins в нескольких проектах, в основном используя Apache Flink. Классический пример — обогащение потока кликов по рекламе (clicks_stream) информацией о соответствующих рекламных кампаниях (campaigns_stream), которая может обновляться.

Основные сложности и подходы к их решению:

  1. Несовпадение временных меток (skew): События из разных потоков приходят с задержкой. Решение — использование Watermarks для определения, когда можно считать окно данных «закрытым» и эмиттить результат.
  2. Хранение состояния: Для join необходимо хранить события из одного потока, ожидая прихода соответствующих событий из другого. В Flink это состояние хранится в managed state (RocksDB). Критично настраивать TTL для этого состояния, чтобы оно не росло бесконечно.
  3. Типы joins: В зависимости от задачи использовал разные стратегии.

Пример реализации Interval Join во Flink (Java API):

DataStream<ClickEvent> clicks = ...;
DataStream<CampaignUpdate> campaigns = ...;

// Присваиваем водяные знаки и ключи
clicks
    .assignTimestampsAndWatermarks(<стратегия>)
    .keyBy(click -> click.getCampaignId())
    .intervalJoin(campaigns
        .assignTimestampsAndWatermarks(<стратегия>)
        .keyBy(campaign -> campaign.getId())
    )
    .between(Time.minutes(-5), Time.minutes(1)) // Клик должен быть в интервале от 5 мин до кампании до 1 мин после
    .process(new ProcessJoinFunction<ClickEvent, CampaignUpdate, EnrichedClick>() {
        @Override
        public void processElement(ClickEvent click, CampaignUpdate campaign, Context ctx, Collector<EnrichedClick> out) {
            out.collect(new EnrichedClick(click, campaign.getDetails()));
        }
    });

Альтернатива в Spark Structured Streaming (Python):

# Определяем водяные знаки для каждого потока
clicks_with_watermark = clicks_df 
    .withWatermark("click_time", "2 minutes") 
    .select("campaign_id", "click_time", "user_id")

campaigns_with_watermark = campaigns_df 
    .withWatermark("update_time", "10 minutes") 
    .select("campaign_id", "update_time", "campaign_name", "budget")

# Выполняем join с условием по времени
joined_stream = clicks_with_watermark.join(
    campaigns_with_watermark,
    expr("""
        campaign_id = campaign_id AND
        click_time >= update_time AND
        click_time <= update_time + interval 1 hour
    """),
    "leftOuter"  # Чтобы получить клики даже без актуальной кампании
)

Ключевые выводы из практики:

  • Мониторинг состояния: Необходимо внимательно следить за размером состояния в операторах join и настраивать соответствующий TTL.
  • Выбор окна: Правильно выбранный интервал join — компромисс между полнотой данных (большое окно) и потреблением памяти/задержкой (малое окно).
  • Тестирование: Такие пайплайны особенно сложно тестировать. Мы использовали подход с «впрыском» тестовых событий с конкретными временными метками в симуляцию runtime Flink.