Ответ
Да, реализовывал stream-stream joins в нескольких проектах, в основном используя Apache Flink. Классический пример — обогащение потока кликов по рекламе (clicks_stream) информацией о соответствующих рекламных кампаниях (campaigns_stream), которая может обновляться.
Основные сложности и подходы к их решению:
- Несовпадение временных меток (skew): События из разных потоков приходят с задержкой. Решение — использование Watermarks для определения, когда можно считать окно данных «закрытым» и эмиттить результат.
- Хранение состояния: Для join необходимо хранить события из одного потока, ожидая прихода соответствующих событий из другого. В Flink это состояние хранится в managed state (RocksDB). Критично настраивать TTL для этого состояния, чтобы оно не росло бесконечно.
- Типы 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.