【发布时间】:2022-04-15 21:10:01
【问题描述】:
关于设置 kafka producer 属性 - enable.idempotence 为 true
kafkaProps.put("enable.idempotence" , "true");
我遇到了错误 -
2021-04-18 16:43:53.584[0;39m [31mERROR[0;39m [35m15524[0;39m [2m---[0;39m [2m[ad | producer-1][0;39m [36mo.a.k.clients.producer.internals.Sender [0;39m [2m:[0;39m [Producer clientId=producer-1] Aborting producer batches due to fatal error
org.apache.kafka.common.errors.ClusterAuthorizationException: Cluster authorization failed.
[2m2021-04-18 16:43:53.585[0;39m [31mERROR[0;39m [35m15524[0;39m [2m---[0;39m [2m[ restartedMain][0;39m [36mc.a.c.g.kafkaclient.PricerProducer [0;39m [2m:[0;39m sending above record failed. java.util.concurrent.ExecutionException: org.apache.kafka.common.errors.ClusterAuthorizationException: Cluster authorization failed.
[2m
集群是否必须支持/启用此功能。如果是这样,集群应该使用的 Kafka 的最低版本是多少。
来自 kafka 文档 -
启用.幂等
当设置为“真”时,生产者将确保只有一份 每条消息都写入流中。如果'false',生产者重试 由于broker故障等原因,可能会写入重复的重试 流中的消息。请注意,启用幂等性需要 max.in.flight.requests.per.connection 小于或等于 5, 重试次数大于 0 且 acks 必须为“全部”。如果这些值 用户未明确设置,将选择合适的值。如果 如果设置了不兼容的值,则会抛出 ConfigException。 类型:布尔默认值:false 有效值:重要性:低
我已设置 - max.in.flight.requests.per.connection=1 并且未设置 acks,因此它自动设置为 -1(全部)。所以,我看到我的配置很好,但即使否则它应该导致 ConfigException 而不是 ClusterAuthorizationException。
【问题讨论】:
-
您使用的是什么版本?而且幂等性和授权无关,所以你是不是想用SASL/SSL连接?
-
我确信我们的经纪人在 0.11 以上。
Kafka 0.11.0 includes support for idempotent and transactional capabilities in the producer.。我们正在使用 SASL/JAAS -
如果我删除
kafkaProps.put("enable.idempotence" , "true");它工作正常。没有抛出 ClusterAuthorizationException。 -
我记得当经纪人上的
log.message.format.version小于 0.11 时,我们会遇到某种类型的错误...还有一个IdempotentWriteACL - docs.confluent.io/5.3.0/kafka/… -
++谢谢。这应该是最可能的原因 -
Enabling Authorization for Idempotent and Transactional APIs
标签: apache-kafka kafka-producer-api