【发布时间】:2019-05-23 13:10:26
【问题描述】:
我正在尝试在 kafka-console-producer 上启用“幂等”选项。 参考以下链接:
- https://gerardnico.com/dit/kafka/producer#idempotent
- https://gerardnico.com/dit/kafka/kafka-console-producer
使用的命令:
$KAFKA_HOME/bin/kafka-console-producer.sh --broker-list node1.com:6667 --topic my_topic --security-protocol SASL_PLAINTEXT --producer-property acks=all --producer-property retries=Integer.MAX_VALUE --producer-property enable.idempotence=true
观察到以下异常:
org.apache.kafka.common.KafkaException: 无法构造 kafka 制片人 在 org.apache.kafka.clients.producer.KafkaProducer.(KafkaProducer.java:433) 在 org.apache.kafka.clients.producer.KafkaProducer.(KafkaProducer.java:291) 在 kafka.producer.NewShinyProducer.(BaseProducer.scala:40) 在 kafka.tools.ConsoleProducer$.main(ConsoleProducer.scala:50) 在 kafka.tools.ConsoleProducer.main(ConsoleProducer.scala) 引起:org.apache.kafka.common.config.ConfigException:必须设置 确认所有人以使用幂等生产者。否则我们 不能保证幂等性。 在 org.apache.kafka.clients.producer.KafkaProducer.configureAcks(KafkaProducer.java:510) 在 org.apache.kafka.clients.producer.KafkaProducer.(KafkaProducer.java:375)
虽然 acks 已设置为“全部”,但我们观察到此异常。 我错过了什么?
以下是使用的版本:
- 经纪人 - 1.0.0
- client - 与 broker 1.0.0 捆绑的控制台生产者
更新
我可以使用回复中建议的--request-required-acks -1 选项在控制台生产者上启用幂等性。
但是,我得到了 ClusterAuthorizationException。
bash$ $KAFKA_HOME/bin/kafka-console-producer.sh --broker-list borker1:6667 --topic my_topic --producer-property enable.idempotence=true --request-required-acks -1 --security-protocol SASL_PLAINTEXT --property "parse.key=true" --property "key.separator=:"
>key1:value1
>[2018-12-26 04:00:56,074] ERROR [Producer clientId=console-producer] Aborting producer batches due to fatal error (org.apache.kafka.clients.producer.internals.Sender)
org.apache.kafka.common.errors.ClusterAuthorizationException: Cluster authorization failed.
[2018-12-26 04:00:56,080] ERROR Error when sending message to topic orm_c1_prv_non_sepa_ci with key: 4 bytes, value: 6 bytes with error: (org.apache.kafka.clients.producer.internals.ErrorLoggingCallback)
org.apache.kafka.common.errors.ClusterAuthorizationException: Cluster authorization failed.
此异常仅在启用幂等选项时发生。没有此选项也可以生成消息。
bash$ $KAFKA_HOME/bin/kafka-console-producer.sh --broker-list broker1:6667 --topic my_topic --security-protocol SASL_PLAINTEXT --property "parse.key=true" --property "key.separator=:"
>key1:value1
>key2:value2
我错过了什么?
【问题讨论】:
-
@cricket_007 已更新相关版本信息
标签: apache-kafka kafka-producer-api