【发布时间】:2021-06-23 05:28:50
【问题描述】:
我有一个简单的 Kafka 流场景,我正在执行 groupyByKey 然后 reduce 然后执行操作。源主题中可能存在重复事件,因此 groupyByKey 和 reduce
该操作可能会出错,在这种情况下,我需要流应用程序来重新处理该事件。在下面的示例中,我总是抛出一个错误来证明这一点。
动作只发生一次且至少发生一次是非常重要的。
我发现的问题是,当流应用程序重新处理事件时,reduce 函数被调用,并且当它返回 null 时,该操作不会被调用。
由于源主题 TOPIC_NAME 只产生一个事件,我希望 reduce 没有任何值并跳到 mapValues。
val topologyBuilder = StreamsBuilder()
topologyBuilder.stream(
TOPIC_NAME,
Consumed.with(Serdes.String(), EventSerde())
)
.groupByKey(Grouped.with(Serdes.String(), EventSerde()))
.reduce { current, _ ->
println("reduce hit")
null
}
.mapValues { _, v ->
println(Id: "${v.correlationId}")
throw Exception("simulate error")
}
为了引起这个问题,我运行了两次流应用程序。这是输出:
首次运行
Id: 90e6aefb-8763-4861-8d82-1304a6b5654e
11:10:52.320 [test-app-dcea4eb1-a58f-4a30-905f-46dad446b31e-StreamThread-1] ERROR org.apache.kafka.streams.KafkaStreams - stream-client [test-app-dcea4eb1-a58f-4a30-905f-46dad446b31e] All stream threads have died. The instance will be in error state and should be closed.
第二次运行
reduce hit
正如您所见,.mapValues 在第二次运行时没有被调用,即使它在第一次运行时出错,导致流应用再次重新处理相同的事件。
是否有可能让流应用程序以减少的步骤重新处理事件,从而以前所未有的方式处理事件? - 或者有更好的方法来解决我的问题吗?
【问题讨论】:
-
不不不,你很困惑。您不需要重新处理事件。总之,忘记一切。你能跟我解释一下,你的业务问题是什么?你的 key-value 输入主题是什么样子的,你期望 POJO 的结果是什么?
-
“您不需要重新处理事件” - 如果事件处理引发异常,Kafka 流会重新处理它,否则您会丢失数据。无论如何,答案是恰好一次语义。
标签: kotlin apache-kafka apache-kafka-streams