【发布时间】: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