【发布时间】:2018-09-27 09:16:58
【问题描述】:
我正在开发一个开发环境,我的系统上有 3 个(dockerized)kafka 代理。 代理的 transaction.state.log.replication.factor 设置为 3。
在流应用程序配置中,我将 processing.guarantee 设置为 EXACTLY_ONCE,在消费者应用程序配置中,我将isolation.level 设置为“read_committed”。
我已经检查了https://docs.confluent.io/current/streams/developer-guide/config-streams.html#processing-guarantee 上的其他配置,并根据指南设置了我的环境。
在从读取状态存储并使用 context.forward(..) 函数生成 100 条消息的流应用程序生成消息一分钟后,消费者应用程序停止读取,就好像在分配的分区上没有任何已提交的消息一样.
一段时间后,流应用程序崩溃并出现以下错误:
“由于致命错误而中止生产者批次 org.apache.kafka.common.errors.ProducerFencedException:生产者 尝试使用旧时代进行操作。要么有更新的 具有相同 transactionalId 的生产者,或生产者的交易 已被代理过期。”
流生产者好像收不到ack,事务过期了。
编辑 1: 当我停止流应用程序时,消费者会收到提交的消息。
【问题讨论】:
-
您是否在任何地方提交事务?显示一些代码示例
-
很难说。我建议检查代理和流日志。
-
找到解决方案了吗?
-
@Rolintocour 看到我的答案:)
标签: java apache-kafka apache-kafka-streams