【问题标题】:Kafka streams: groupByKey and reduce not triggering action exactly once when error occurs in streamKafka流:groupByKey和reduce在流中发生错误时不会恰好触发一次动作
【发布时间】:2021-06-23 05:28:50
【问题描述】:

我有一个简单的 Kafka 流场景,我正在执行 groupyByKey 然后 reduce 然后执行操作。源主题中可能存在重复事件,因此 groupyByKeyreduce 该操作可能会出错,在这种情况下,我需要流应用程序来重新处理该事件。在下面的示例中,我总是抛出一个错误来证明这一点。

动作只发生一次且至少发生一次是非常重要的。

我发现的问题是,当流应用程序重新处理事件时,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


【解决方案1】:

我缺少流应用程序的属性设置。

props["processing.guarantee"]= "exactly_once"

通过设置此项,它将保证从拾取事件点创建的任何状态都将回滚,以防引发异常并且流应用程序崩溃。

问题在于流应用程序会再次获取事件以重新处理,但减速器步骤的状态一直存在。通过启用exactly_once 设置,它可以确保reducer 状态也被回滚。

它现在成功地重新处理该事件,就好像它以前从未见过它一样

【讨论】:

    猜你喜欢
    • 2015-09-11
    • 2015-05-17
    • 1970-01-01
    • 2020-06-28
    • 2022-06-22
    • 2020-05-12
    • 1970-01-01
    • 2019-07-29
    • 1970-01-01
    相关资源
    最近更新 更多