【问题标题】:Kafka Wheel Timer卡夫卡轮式计时器
【发布时间】:2019-07-13 02:09:03
【问题描述】:

美好的一天,

我想知道kafka队列是否可以保持数据几秒钟而不是释放数据。

我收到来自 kafka 主题的消息, 解析数据后,我将其在内存中保存一段时间(10 秒)(这会随着唯一消息的传递而增加),每条消息都有自己的计时器),我希望 kafka 告诉我该消息已过期(10秒),以便我可以继续其他任务。

但由于 flink/kafka 是事件驱动的,我希望 kafka 有某种圆形计时轮,可以在 10 秒后将消息的密钥重现给消费者。

知道如何使用 flink 窗口或 kafka 功能来归档它吗?

问候

【问题讨论】:

    标签: scala timer apache-kafka apache-flink


    【解决方案1】:

    关于你最初的问题:

    我想知道kafka队列是否可以保持数据几秒钟而不是释放数据

    您可以将log.cleanup.policy 设置为delete(这是默认设置)并将retention.ms 从默认604800000(1 周)更改为10000

    您能否再次解释一下您还想检查什么,Regards 部分之后的意思是什么?

    【讨论】:

    • 我认为 hold Ideas frontier 意味着推迟转发数据,而不是将它们从代理中删除。
    • 是的,这正是我需要的@wardziniak
    【解决方案2】:

    您可以更深入地了解 Kafka Streams 库。 https://kafka.apache.org/21/documentation/streams/developer-guide/dsl-api.html, https://kafka.apache.org/21/documentation/streams/developer-guide/processor-api.html.

    使用 Kafka Streams,您可以完成许多复杂的事件处理工作。处理器 API 是较低级别的 API,为您提供更大的灵活性,例如每个处理消息放入状态存储(Kafka Streams 抽象,复制到更改日志主题),然后使用Punctuator,您可以检查消息是否过期 em>。

    【讨论】:

    • 感谢您研究这些建议
    猜你喜欢
    • 2013-11-04
    • 1970-01-01
    • 2020-08-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-09-03
    • 1970-01-01
    • 2018-08-29
    相关资源
    最近更新 更多