【问题标题】:Failure to Produce to Embedded Kafka Broker无法生成到嵌入式 Kafka 代理
【发布时间】:2017-05-10 03:04:18
【问题描述】:

我正在尝试使用 kafka 的 kafka.zk.EmbeddedZookeeperkafka.server.KafkaServerreturned by kafka.utils.TestUtils/createServer 来运行 kafka 服务器进行测试。

但我遇到了一个障碍,即尝试发送消息超时,并返回 KafkaProducer$Future 失败。下面是我正在使用的 kafka 版本。下面的代码是 Clojure 与 Kafka 库的互操作。

[org.apache.kafka/kafka_2.11 "0.10.0.1"]
[org.apache.kafka/kafka-clients "0.10.1.0"]

这就是我能走多远。

  • Zookeeper 端口是随机分配的(请参阅here)。
  • 可以成功创建 Zookeeper server 并使用netcat 连接到它。
  • 可以成功创建主题。
  • 可以使用 netcat 成功创建 Kafka 代理并连接到它。
  • 第 5 步是流程失败的地方。

这个SO question 表明传入正确的Time 对象很重要。但是MockTime 看起来是一个合理的实现。以前有人解决过这个问题吗?

;; 1. Create Zookeeper
(require '[clojure.test :refer :all]
     '[kafkaesque.topics :as kt]
     '[kafkaesque.utils :as ku]
     '[clojure.pprint :refer [pprint]])

(import '[java.nio.file Files]
    '[kafka.zk EmbeddedZookeeper]
    '[kafka.server KafkaServer KafkaConfig]
    '[kafka.utils TestUtils Time MockTime])

(def zk-config {:zkhost "127.0.0.1"})
(def topic-name "client-test")
(def ^EmbeddedZookeeper zkServer (EmbeddedZookeeper.))


;; 2. Create Kafka Broker
(def zk-connect-str (str "127.0.0.1" ":" (.port zkServer)))
(def zku ((ZkUtils/apply (ZkUtils/createZkClient zk-connect-str 10000 8000) false)))
(def brokerhost "127.0.0.1")
(def brokerport "9092")

(def ^KafkaConfig config (KafkaConfig. {"zookeeper.connect" zk-connect-str
                       "broker.id" "0"
                       "log.dirs" (.toString
                           (.toAbsolutePath
                            (Files/createTempDirectory
                             "kafka-" (make-array java.nio.file.attribute.FileAttribute 0))))
                       "listeners" (str "PLAINTEXT://"  brokerhost  ":"  brokerport)}))

(def ^Time mock (MockTime.))
(def ^KafkaServer kafkaServer (TestUtils/createServer config mock))


;; 3. Create a Topic
(kt/create! zku topic-name 1 1 {})
(kt/topic-exists? zku topic-name)   ;; returns true


;; 4. Create a Producer and ProducerRecord
(def producer-a (kc/producer {"bootstrap.servers" "127.0.0.1:9092"
                 "acks"              "all"
                 "retries"           "0"
                 "batch.size"        "16384"
                 "linger.ms"         "1"
                 "buffer.memory"     "33554432"
                 "key.serializer"    "org.apache.kafka.common.serialization.StringSerializer"
                 "value.serializer"  "org.apache.kafka.common.serialization.StringSerializer"}))

(def message-key "k1")
(def message-value "foobar")
(def record-a (kc/producer-record topic-name 0 message-key message-value))


;; 5. Send a message
(def send-result (kc/send! producer-a record-a))  ;; Times out, and returns a KafkaProducer$Future failure.

【问题讨论】:

  • 可以不为 TestUtils.createServer 指定 Time 实例,因为它将创建一个默认实例。至于“超时”,您能否确认它是说获取元数据超时?此外,我注意到客户端和服务器的版本不匹配。可以用同版本的Kafka重试吗?
  • @amethystic Crikey,就是这样 - 版本不匹配。 org.apache.kafka/kafka_2.11 "0.10.0.1",对比 org.apache.kafka/kafka-clients " 0.10.1.0"。哎呀,我把头发拉出来了!我将所有版本移至 "0.10.1.0",一切正常。干杯:)

标签: scala clojure apache-kafka apache-zookeeper


【解决方案1】:

感谢 @amethystic 指出我的阅读障碍 :) 问题是版本不匹配。我正在使用 org.apache.kafka/kafka_2.11 "0.10.0.1",与 org.apache.kafka/kafka-clients “0.10.1.0”。

我将所有版本移至 “0.10.1.0”,一切正常。

希望这会有所帮助。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-08-31
    • 1970-01-01
    • 2018-06-16
    • 2018-10-18
    • 2021-12-22
    • 1970-01-01
    相关资源
    最近更新 更多