You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在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_login Kafka主题,供下游消费者使用。

现有尝试与顾虑

  • 最初考虑使用时长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;

核心需求

  1. 缓存Kafka主题数据X小时
  2. X小时后自动过期并回收缓存数据
  3. 避免生成过多缓存副本或重复执行关联:理想状态下每个缓存元素仅保留一份,仅保留故障转移所需的副本
  4. 尽量用Beam SQL实现核心逻辑,减少Java代码编写量

内容的提问来源于stack exchange,提问作者Gergely

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 20:02:34