【问题标题】:Publishing a message to Kafka running inside docker向在 docker 内运行的 Kafka 发布消息
【发布时间】:2016-09-14 08:57:13
【问题描述】:

我在 docker 容器中运行 Kafka。我使用以下命令启动我的容器

docker run --rm -p 2181:2181 -p 9092:9092 -p 8081:8081 --env 
ADVERTISED_HOST=\`docker-machine ip \\`docker-machine active\\`` --env 
ADVERTISED_PORT=9092 -v  
/Users/abhishek.srivastava/MyProjects/KafkaTest/target/scala-2.11:/app 
-it --  name kafka spotify/kafka bash

我写了一个简单的程序,我可以在容器中复制并执行它,它运行良好。

object KafkaProducerString {

  def SendStringMessage(msg: String) : Unit = {
    val inputRecord = new ProducerRecord[String, String]("test", null, msg)
    val producer: KafkaProducer[String, String] = CreateProducerString
    val rm = producer.send(inputRecord).get(10, SECONDS)
    println(s"offset: ${rm.offset()} partition: ${rm.partition()} topic: ${rm.topic()}")
    producer.close()
  }

  private def CreateProducerString: KafkaProducer[String, String] = {
    val props = new Properties()
    props.put("bootstrap.servers", "localhost:9092")
    props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
    props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
    props.put("batch.size", "0")
    props.put("client.id", "1")
    val producer = new KafkaProducer[String, String](props)
    producer
  }
}

但是,如果我从容器外部(从我的 mac)运行相同的程序。 [我将“localhost”替换为docker-machine ip的输出]

我收到此错误

[error] (run-main-0) java.util.concurrent.TimeoutException: Timeout after waiting for 10000 ms.
java.util.concurrent.TimeoutException: Timeout after waiting for 10000 ms.
    at org.apache.kafka.clients.producer.internals.FutureRecordMetadata.get(FutureRecordMetadata.java:50)
    at org.apache.kafka.clients.producer.internals.FutureRecordMetadata.get(FutureRecordMetadata.java:25)
    at com.abhi.KafkaProducerString$.SendStringMessage(KafkaProducerString.scala:23)
    at com.abhi.KafkaMain$$anonfun$main$1.apply$mcVI$sp(KafkaMain.scala:19)
    at scala.collection.immutable.Range.foreach$mVc$sp(Range.scala:160)
    at com.abhi.KafkaMain$.main(KafkaMain.scala:17)
    at com.abhi.KafkaMain.main(KafkaMain.scala)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:57)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:606)

我的理解是,对于远程的 kafka 生产者,我需要打开的唯一端口是 2181(zookeeper)和 9092(kafka),你可以看到我已经打开了这些。

但同样的程序在容器外执行时仍然超时,但在容器内(使用本地主机)时可以工作。

编辑::根据下面的建议,我尝试了以下

docker run --rm -p 127.0.0.1:2181:2181 -p 127.0.0.1:9092:9092 -p 
127.0.0.1:8081:8081 --env ADVERTISED_HOST=`docker-machine ip \`docker-machine 
active\`` --env ADVERTISED_PORT=9092 -v 
/Users/abhishek.srivastava/MyProjects/KafkaTest/target/scala-2.11:/app -it --
name kafka kafka_9.0 bash

和

docker run --rm -p 0.0.0.0:2181:2181 -p 0.0.0.0:9092:9092 -p 0.0.0.0:8081:8081 
--env ADVERTISED_HOST=`docker-machine ip \`docker-machine active\`` --env 
ADVERTISED_PORT=9092 -v 
/Users/abhishek.srivastava/MyProjects/KafkaTest/target/scala-2.11:/app -it --
name kafka kafka_9.0 bash

但是这些并没有解决问题。我遇到了完全相同的问题

【问题讨论】:

  • 嗨,我也面临同样的问题。你解决了吗?
  • 我放弃了 :) 有时间会再试一次。检查下面的解决方案。让我知道它是否有效:)
  • 我正在使用 wurstmeister kafka docker。一切都在 docker 中运行,但我在主机上的生产者/消费者代码无法连接到 kafka 代理。我正在解决这个问题,如果出现问题我会告诉你
  • 即使我也面临同样的问题。如果您找到任何解决方案,请更新。或者是否有任何适合您的替代方法。

标签: docker apache-kafka kafka-producer-api


【解决方案1】:

您必须将 docker 容器绑定到本地计算机。这可以通过使用 docker run as 来完成:

docker run --rm -p 127.0.0.1:2181:2181 -p 127.0.0.1:9092:9092 -p 127.0.0.1:8081:8081 ....

或者,您可以使用绑定 IP 的 docker run:

docker run --rm -p 0.0.0.0:2181:2181 -p 0.0.0.0:9092:9092 -p 0.0.0.0:8081:8081 .....

如果您想让 docker 容器在您的网络上可路由,您可以使用:

docker run --rm -p <private-IP>:2181:2181 -p <private-IP>:9092:9092 -p <private-IP>:8081:8081 ....

或者最后,您可以通过以下方式不将您的网络接口容器化:

docker run --rm -p 2181:2181 -p 9092:9092 -p 8081:8081 --net host ....

【讨论】:

  • 我尝试了前两个,但没有解决我的问题
  • 您是否在 Dockerfile 中公开了相关端口?
  • 您必须首先公开您尝试绑定的端口 8081。还要检查 docker 容器是否正在使用 docker ps 运行。并确保如果 docker 容器正在运行,您可以在主机上绑定到主机的暴露端口上远程登录。
  • 试过了......同样的问题。我想知道您是否可以从 github 下载我的项目,然后尝试查看是否可以运行它(从 docker 容器外部)github.com/abhitechdojo/KafkaTest
【解决方案2】:

虽然我自己也面临类似的问题,但我可以尝试解释这种行为。

Kafka 生产者会在将记录发布到 Topic 之前从 Zookeeper 中查找分区领导者。 Zookeeper 将拥有由 Kafka 服务器标记的领导主机条目,该服务器在 Docker 容器内运行。

因此,服务器标记的 IP 将是 Docker 内部 IP,而不是主机 IP。这当然不能从客户端机器上解决,因此会超时。

可能的解决方案是将advertised.host.name 设置为Docker 机器的主机IP。但是,这会引入另一个问题(正如我所面临的那样!)

服务器获取代理元数据现在将开始失败。这是因为现在 Zookeeper 条目具有主机 IP,无法从容器内部解析。因此,任何消费者应用程序现在都会开始收到LEADER_NOT_AVAILABLE 警告。

这是一种死锁情况,解决方案主要取决于所采用的主机解析策略。我想知道人们会建议如何去这里。

编辑:最后我们使用主机网络 [--net=host] 并使用节点静态 IP 来解决问题。

【讨论】:

  • 你可能说了我遇到的一个问题
猜你喜欢
  • 1970-01-01
  • 2017-05-14
  • 1970-01-01
  • 2016-12-07
  • 1970-01-01
  • 1970-01-01
  • 2019-05-11
  • 2023-03-17
相关资源
最近更新 更多