【问题标题】:Apache Flink delay processing of certain eventsApache Flink 延迟处理某些事件
【发布时间】:2021-10-05 06:40:12
【问题描述】:

我需要延迟处理某些事件。

例如。我有三个事件(在 Kafka 上发布):

  • A (id: 1, retryAt: now)
  • B(id:2,retryAt:10 分钟后)
  • C (id: 3, retryAt: now)

我需要立即处理记录 A 和 C,而记录 B 需要十分钟后处理。 这在 Apache Flink 中是否可行?

到目前为止,无论我研究过什么,似乎“触发器”可能有助于在 Flink 中实现它,但还不能正确实现它。

我也查看了 Kafka 文档,但在那里看起来并不可行。

【问题讨论】:

    标签: apache-kafka apache-flink flink-streaming flink-sql flink-cep


    【解决方案1】:

    触发器适用于窗口,但窗口化似乎不适合您的用例。

    更好的解决方案是使用带有KeyedProcessFunction 的定时器。根据您是要等待 10 分钟的处理时间还是 10 分钟的事件时间,您将选择处理时间计时器或事件时间计时器。

    您还需要使用 Flink 状态来存储需要稍后处理的事件。

    您将找到流程函数here 的文档。 Flink 训练中还有一些额外的例子,herehere

    FWIW,Flink 的 Stateful Functions API 可能更适合您的工作,在这种情况下您可以使用 delayed messages

    【讨论】:

    • 谢谢大卫。这是有用的信息。 “Flink 状态来存储稍后处理的事件”:为此,我需要使用 RocksDB 还是任何外部 DB。因为,当 Flink 节点宕机时,In-Memory 存储会丢失消息。
    • Flink 依靠检查点来实现容错。工作状态由状态后端管理,如果选择HashMapStateBackend,则存储在内存中,如果选择EmbeddedRocksDBStateBackend,则存储在本地磁盘上。
    • 嗨,David,您能就另一个 Flink 问题发表您的看法吗? stackoverflow.com/questions/68725502/…
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-11-19
    相关资源
    最近更新 更多