【问题标题】:UnknownProducerIdException in Kafka streams when enabling exactly once仅启用一次时,Kafka 流中的 UnknownProducerIdException
【发布时间】:2018-09-27 02:37:06
【问题描述】:

在 Kafka 流应用程序上启用恰好一次处理后,日志中出现以下错误:

ERROR o.a.k.s.p.internals.StreamTask - task [0_0] Failed to close producer 
due to the following error:

org.apache.kafka.streams.errors.StreamsException: task [0_0] Abort 
sending since an error caught with a previous record (key 222222 value 
some-value timestamp 1519200902670) to topic exactly-once-test-topic- 
v2 due to This exception is raised by the broker if it could not 
locate the producer metadata associated with the producerId in 
question. This could happen if, for instance, the producer's records 
were deleted because their retention time had elapsed. Once the last 
records of the producerId are removed, the producer's metadata is 
removed from the broker, and future appends by the producer will 
return this exception.
  at org.apache.kafka.streams.processor.internals.RecordCollectorImpl.recordSendError(RecordCollectorImpl.java:125)
  at org.apache.kafka.streams.processor.internals.RecordCollectorImpl.access$500(RecordCollectorImpl.java:48)
  at org.apache.kafka.streams.processor.internals.RecordCollectorImpl$1.onCompletion(RecordCollectorImpl.java:180)
  at org.apache.kafka.clients.producer.KafkaProducer$InterceptorCallback.onCompletion(KafkaProducer.java:1199)
  at org.apache.kafka.clients.producer.internals.ProducerBatch.completeFutureAndFireCallbacks(ProducerBatch.java:204)
  at org.apache.kafka.clients.producer.internals.ProducerBatch.done(ProducerBatch.java:187)
  at org.apache.kafka.clients.producer.internals.Sender.failBatch(Sender.java:627)
  at org.apache.kafka.clients.producer.internals.Sender.failBatch(Sender.java:596)
  at org.apache.kafka.clients.producer.internals.Sender.completeBatch(Sender.java:557)
  at org.apache.kafka.clients.producer.internals.Sender.handleProduceResponse(Sender.java:481)
  at org.apache.kafka.clients.producer.internals.Sender.access$100(Sender.java:74)
  at org.apache.kafka.clients.producer.internals.Sender$1.onComplete(Sender.java:692)
  at org.apache.kafka.clients.ClientResponse.onComplete(ClientResponse.java:101)
  at org.apache.kafka.clients.NetworkClient.completeResponses(NetworkClient.java:482)
  at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:474)
  at org.apache.kafka.clients.producer.internals.Sender.run(Sender.java:239)
  at org.apache.kafka.clients.producer.internals.Sender.run(Sender.java:163)
  at java.lang.Thread.run(Thread.java:748)
Caused by: org.apache.kafka.common.errors.UnknownProducerIdException

我们用一个最小的测试用例重现了这个问题,我们将消息从一个源流移动到另一个流而不进行任何转换。源流包含数月产生的数百万条消息。 KafkaStreams 对象是使用以下 StreamsConfig 创建的:

  • StreamsConfig.PROCESSING_GUARANTEE_CONFIG = "exactly_once"
  • StreamsConfig.APPLICATION_ID_CONFIG = "一些应用 ID"
  • StreamsConfig.NUM_STREAM_THREADS_CONFIG = 1
  • ProducerConfig.BATCH_SIZE_CONFIG = 102400

应用程序能够在异常发生之前处理一些消息。

上下文信息:

  • 我们正在运行一个包含 5 个 Zookeeper 节点的 5 个节点的 Kafka 1.1.0 集群。
  • 有多个应用实例正在运行

有没有人以前见过这个问题,或者可以给我们任何关于可能导致这种行为的提示?

更新

我们从头开始创建了一个新的 1.1.0 集群,并开始毫无问题地处理新消息。但是,当我们从旧集群中导入旧消息时,我们会在一段时间后遇到相同的 UnknownProducerIdException。

接下来,我们尝试将接收器主题上的cleanup.policy 设置为compact,同时将retention.ms 保持3 年。现在错误没有发生。但是,消息似乎已丢失。源偏移量1.06亿,汇偏移量1亿。

【问题讨论】:

  • 生产者 ID 直接存储在日志中——因此,如果您的所有数据都被删除,则 PID 可能会丢失。您的保留时间是多少?您的数据有哪些时间戳?
  • 源主题和目标主题的保留时间设置为3年。生产者 ID 存储在日志中是什么意思?
  • 具体来说,我们将retention.ms设置为3年,将cleanup.policy设置为删除
  • 消息时间戳是否早于保留时间?写入生产者的生产者 ID 存储为每条消息,日志用作已知 PID 的“事实来源”(PID 不会再次存储在其他地方)。也可能是生产者经纪人的错误......不确定。

标签: apache-kafka apache-kafka-streams


【解决方案1】:

正如 cmets 中所解释的,目前似乎存在一个错误,当重播超过(最大可配置?)保留时间的消息时可能会导致问题。

在撰写本文时,此问题尚未解决,始终可以在此处查看最新状态:

https://issues.apache.org/jira/browse/KAFKA-6817

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-04-18
    • 2020-09-15
    • 2018-10-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多