【问题标题】:How to enable idempotence in kafka-console-producer?如何在 kafka-console-producer 中启用幂等性?
【发布时间】:2019-05-23 13:10:26
【问题描述】:

我正在尝试在 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


【解决方案1】:

您不能为 ConsoleProducer 设置 acksproducer-property。改用request-required-acks,如下图:

bin/kafka-console-producer.sh --broker-list localhost:9092 --topic test --producer-property enable.idempotence=true --request-required-acks -1

【讨论】:

    猜你喜欢
    • 2018-03-23
    • 2018-10-03
    • 2021-04-26
    • 2016-12-06
    • 1970-01-01
    • 2021-11-25
    • 1970-01-01
    • 1970-01-01
    • 2017-12-01
    相关资源
    最近更新 更多