【问题标题】:Event retry mechanism in Apache FlinkApache Flink 中的事件重试机制
【发布时间】:2021-07-18 07:57:59
【问题描述】:

我正在从 Kafka 读取多个主题,然后进行状态操作,然后再次保存到 Kafka。这是流程;

val stream1 = topic1.map(prepareData())
val stream2 = topic2.map(prepareData1())
val delayStream = delayqueue.map(prepareData2()) // these messages comes from delay

val enrichedResults = stream1.union(stream2).union(delayStream).map(stateOperation)
val taggedStream = enrichedResults.process(tagStreamAsRetryOrNot)
val retryMessages = taggedStream.getSideOutput(tag)

retryMessages.addSink(kafkaDelayQueue1)
taggedStream.addSink(targetTopic)

在这个流程中,如果找不到任何状态,这个消息将被发送到delay queue。稍后我会处理此消息。

延迟流程的另一面:

delayQueue(written by flink) => consumerAPP(check retrycounts and timestamps) => anotherQueue

(这会被flink再次消费。我不想再次向主主题发送消息,因为主主题也被其他团队消费。)

这种方法好吗?可以用更少的努力改善这种流程吗?或者有什么最佳实践吗?

【问题讨论】:

  • 为什么需要重新处理/重试?也许这可以通过在两个流之间进行某种时间连接来避免。
  • 我使用从 topic1 消息中保存的状态来丰富 topic2 消息。来自 topic2 的消息可能比 topic1 更早。因此,我必须重新处理没有状态的 topic2 事件。

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


【解决方案1】:

通常处理这种情况的方式是在状态中缓冲早期消息,直到携带丰富所需数据的预期消息到达另一个流。例如,您可以在 KeyedCoProcessFunction 中执行此操作。

Flink SQL(和 Table API)的设置使这种流式连接变得容易。

有关该主题的更多信息,请参阅the documentation on SQL joins

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-12-10
    • 2021-10-05
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多