在Apache Beam流处理管道中实现带TTL的缓存数据关联
问题:基于Apache Beam/Dataflow实现带过期缓存的Kafka流关联处理
背景
我正尝试将一批基于命令式代码、依赖Redis缓存的Kafka处理管道,迁移至运行在Dataflow上的Apache Beam。所有管道的输入均来自Kafka主题,分为两类:
- 缓存类:数据需缓存X小时后自动过期
- 非缓存类:数据到达后立即处理
Kafka消费者会通过内连接或外连接,用缓存数据丰富非缓存输入,待处理主题数量在3-12个之间。
示例场景
Kafka输入包含三个按city键分区的主题:
login_event:非缓存类,数据到达立即处理trending_news:缓存类,需缓存12小时announcements:缓存类,需缓存72小时
需要将这三类数据按city做三方关联,把聚合后的富化对象写入enriched_loginKafka主题,供下游消费者使用。
现有尝试与顾虑
- 最初考虑使用时长72小时、周期1小时的滑动窗口,但担心会生成多达72份
announcements缓存副本,且关联操作需在所有"活跃"窗口执行,效率低下。 - 研究过用有状态DoFn和Timers实现带过期时间的缓存,但找不到让缓存状态供Beam SQL查询使用的方法,暂时搁置。
- 了解到ksqldb的流处理可实现类似需求,执行以下SQL时会隐式创建按键分区的临时缓冲区,过期时间设为窗口结束时间,这正是我希望在Beam中实现的行为:
CREATE STREAM enriched_login WITH (kafka_topic=...) AS SELECT le.login_details AS login_details, le.key AS login_key, le.city AS city, tn.news_topics AS news_topics, an.messages AS messages FROM login_event le INNER JOIN trending_news tn WITHIN 12 HOURS GRACE PERIOD 10 MINUTES ON le.city = tn.city INNER JOIN announcements an WITHIN 72 HOURS GRACE PERIOD 10 MINUTES ON an.city = le.city;
核心需求
- 缓存Kafka主题数据X小时
- X小时后自动过期并回收缓存数据
- 避免生成过多缓存副本或重复执行关联:理想状态下每个缓存元素仅保留一份,仅保留故障转移所需的副本
- 尽量用Beam SQL实现核心逻辑,减少Java代码编写量
内容的提问来源于stack exchange,提问作者Gergely
相关产品推荐
相关产品推荐

