【问题标题】:Docker Kafka w/ Python consumer带有 Python 消费者的 Docker Kafka
【发布时间】:2019-02-25 13:46:40
【问题描述】:

我正在使用 dockerized Kafka 并编写了一个 Kafka 消费者程序。当我在本地机器上的 docker 和应用程序中运行 Kafka 时,它可以完美运行。但是当我在 docker 中配置本地应用程序时,我遇到了问题。该问题可能是由于在应用程序启动之前未创建主题。

docker-compose.yml

version: '3'
services:
  zookeeper:
    image: wurstmeister/zookeeper
    ports:
      - "2181:2181"
  kafka:
    image: wurstmeister/kafka
    ports:
      - "9092:9092"
    environment:
      KAFKA_ADVERTISED_HOST_NAME: localhost
      KAFKA_CREATE_TOPICS: "test:1:1"
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
    volumes:
      - /var/run/docker.sock:/var/run/docker.sock
  parse-engine:
    build: .
    depends_on:
      - "kafka"
    command: python parse-engine.py
    ports:
     - "5000:5000"

解析引擎.py

from kafka import KafkaConsumer
import json

try:
    print('Welcome to parse engine')
    consumer = KafkaConsumer('test', bootstrap_servers='localhost:9092')
    for message in consumer:
        print(message)
except Exception as e:
    print(e)
    # Logs the error appropriately. 
    pass

错误日志

kafka_1         | [2018-09-21 06:27:17,400] INFO [SocketServer brokerId=1001] Started processors for 1 acceptors (kafka.network.SocketServer)
kafka_1         | [2018-09-21 06:27:17,404] INFO Kafka version : 2.0.0 (org.apache.kafka.common.utils.AppInfoParser)
kafka_1         | [2018-09-21 06:27:17,404] INFO Kafka commitId : 3402a8361b734732 (org.apache.kafka.common.utils.AppInfoParser)
kafka_1         | [2018-09-21 06:27:17,431] INFO [KafkaServer id=1001] started (kafka.server.KafkaServer)
**parse-engine_1  | Welcome to parse engine
parse-engine_1  | NoBrokersAvailable 
parseengine_parse-engine_1 exited with code 0**
kafka_1         | creating topics: test:1:1

由于我已经在 docker-compose 中添加了 depends_on 属性,但在启动主题应用程序连接之前发生了错误。

我读到我可以在 docker-compose 文件中添加脚本,但我正在寻找一些简单的方法。

感谢您的帮助

【问题讨论】:

  • 不,不一样。我能够连接 Kafka,但面临诸如懒惰主题创建之类的问题。
  • 绝对一样。 NoBrokersAvailable 因为你连接到错误的引导服务器
  • 仍面临问题,因此未标记为已接受
  • 请用您当前的代码更新您的问题,然后,因为如上所述,Python 代码中的localhost:9092 不正确

标签: python docker apache-kafka docker-compose kafka-consumer-api


【解决方案1】:

你的问题是网络。在您的 Kafka 配置中,您正在设置

KAFKA_ADVERTISED_HOST_NAME: localhost

但这意味着任何客户端(包括您的 python 应用程序)都将连接到代理,然后由代理告知使用localhost 进行任何连接。由于来自您的客户端计算机(例如您的 python 容器)的 localhost 不是代理所在的位置,因此请求将失败。

您可以在此处详细阅读有关解决 Kafka 侦听器问题的更多信息

因此,要解决您的问题,您可以执行以下两项操作之一:

  1. 只需更改您的 compose 以使用 Kafka 的 internal 主机名 (KAFKA_ADVERTISED_HOST_NAME: kafka)。这意味着 docker 网络的任何客户端都可以正常访问它,但没有外部客户端能够(例如从您的主机):

     version: '3'
     services:
     zookeeper:
         image: wurstmeister/zookeeper
         ports:
         - "2181:2181"
     kafka:
         image: wurstmeister/kafka
         ports:
         - "9092:9092"
         environment:
         KAFKA_ADVERTISED_HOST_NAME: kafka
         KAFKA_CREATE_TOPICS: "test:1:1"
         KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
         volumes:
         - /var/run/docker.sock:/var/run/docker.sock
     parse-engine:
         build: .
         depends_on:
         - "kafka"
         command: python parse-engine.py
         ports:
         - "5000:5000"
    

    然后您的客户将在 kafka:9092 访问代理,因此您的 python 应用程序将更改为

     consumer = KafkaConsumer('test', bootstrap_servers='kafka:9092')
    
  2. 添加一个新的 Kafka 监听器。这使得它可以在 docker 网络内部和外部访问。端口 29092 用于访问 external 到 docker 网络(例如,从您的主机),9092 用于 internal 访问。

    您仍然需要更改您的 python 程序才能在正确的地址访问 Kafka。在这种情况下,由于它是 Docker 网络的内部,您可以使用:

     consumer = KafkaConsumer('test', bootstrap_servers='kafka:9092')
    

    由于我不熟悉 wurstmeister 图像,所以这个 docker-compose 是基于我知道的 Confluent 图像:

    (编辑把我的yaml弄坏了,你可以找到it here

     ---
     version: '2'
     services:
       zookeeper:
         image: confluentinc/cp-zookeeper:latest
         environment:
           ZOOKEEPER_CLIENT_PORT: 2181
           ZOOKEEPER_TICK_TIME: 2000
    
       kafka:
         # "`-._,-'"`-._,-'"`-._,-'"`-._,-'"`-._,-'"`-._,-'"`-._,-'"`-._,-'"`-._,-
         # An important note about accessing Kafka from clients on other machines: 
         # -----------------------------------------------------------------------
         #
         # The config used here exposes port 29092 for _external_ connections to the broker
         # i.e. those from _outside_ the docker network. This could be from the host machine
         # running docker, or maybe further afield if you've got a more complicated setup. 
         # If the latter is true, you will need to change the value 'localhost' in 
         # KAFKA_ADVERTISED_LISTENERS to one that is resolvable to the docker host from those 
         # remote clients
         #
         # For connections _internal_ to the docker network, such as from other services
         # and components, use kafka:9092.
         #
         # See https://rmoff.net/2018/08/02/kafka-listeners-explained/ for details
         # "`-._,-'"`-._,-'"`-._,-'"`-._,-'"`-._,-'"`-._,-'"`-._,-'"`-._,-'"`-._,-
         #
         image: confluentinc/cp-kafka:latest
         depends_on:
           - zookeeper
         ports:
           - 29092:29092
         environment:
           KAFKA_BROKER_ID: 1
           KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
           KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092
           KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
           KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
           KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
    

免责声明:我为 Confluent 工作

【讨论】:

  • 在尝试第一个解决方案后出现错误“KafkaUnavailableError: All servers failed to process request: [('kafka', 9092, )]”在应用程序中
  • 链接中提供的 yaml:gist.github.com/rmoff/fb7c39cc189fc6082a5fbd390ec92b3d 有错字,对我不起作用
  • @joe 错字是什么?如果你让我知道,我会更新要点
  • 使用 docker 主机的 linux 主机名(即 bash 中的 $(hostname) 的值)而不是 localhost 来启用来自 docker 以外的其他机器的客户端不是更好吗主持人? (如您的参考博文中所示)
  • @RobinMoffatt 如果我从不同的网络访问 kafka 服务器,我是否必须更改 "KAFKA_ADVERTISED_HOST_NAME: PLAINTEXT://ipaddress:9092, BROKER://127.0.0.1:9092" ?这样外部应用程序将被重定向到服务器 IP 地址?
【解决方案2】:

这一行

KAFKA_ADVERTISED_HOST_NAME: localhost

说代理将自己宣传为仅在 localhost 上可用,这意味着所有 Kafka 客户端只会返回自己,而不是实际代理地址的实际列表。如果您的客户端仅位于您的主机上,这将很好 - 请求总是转到本地主机,转发到容器

但是,对于其他容器中的应用程序,它们需要指向 Kafka 容器,所以应该说KAFKA_ADVERTISED_HOST_NAME: kafka,其中kafka 是 Docker Compose 服务的名称。然后其他容器中的客户端会尝试连接到那个容器


话虽如此,那么这一行

consumer = KafkaConsumer('test', bootstrap_servers='localhost:9092')

您将 Python 容器 指向自身,而不是 kafka 容器。

应该改为kafka:9092

【讨论】:

  • 如果我从不同的网络访问 kafka 服务器,是否必须更改 "KAFKA_ADVERTISED_HOST_NAME: PLAINTEXT://ipaddress:9092, BROKER://127.0.0.1:9092" ?这样外部应用程序将被重定向到服务器 IP 地址?
  • @Danielle 没错,除了不使用主机名配置,而是在广告多个地址时广告监听器
  • @OneCricketeer,你说的“......不使用主机名配置,但广告的听众......”是什么意思?你的意思是我们应该使用 IP 地址而不是主机名?
  • 啊,好吧...我认为您的意思是 Danielle 应该在这里使用“KAFKA_ADVERTISED_LISTENERS”而不是“KAFKA_ADVERTISED_HOST_NAME”,因为他使用的值不仅包括主机名,还包括端口号。
  • @IwanSatria 这是正确的 - IMO,总是使用广告听众。值得一提的是,广告主机名(和端口)属性本身已被弃用,取而代之的是 advertised.listeners 代理配置
【解决方案3】:

在我的情况下,我想从本地运行的外部 python 客户端(作为生产者)访问 Kafka 容器,这里是容器和对我有用的 python 代码的组合(平台 MAC OS 和 docker 版本 2.4.0) :

动物园管理员容器:

docker run -d \
-p 2181:2181 \
--name=zookeeper \
-e ZOOKEEPER_CLIENT_PORT=2181 \
confluentinc/cp-zookeeper:5.2.3

kafka 容器:

docker run -d \ 
-p 29092:29092 \ 
-p 9092:9092 \ 
--name=kafka \ 
-e KAFKA_ZOOKEEPER_CONNECT=host.docker.internal:2181 \ 
-e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=BROKER:PLAINTEXT,PLAINTEXT:PLAINTEXT \ 
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:29092,BROKER://localhost:9092 \ 
-e KAFKA_INTER_BROKER_LISTENER_NAME=BROKER \ 
-e KAFKA_BROKER_ID=1 \ 
-e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \ 
-e KAFKA_CREATE_TOPICS="test:1:1" \ 
confluentinc/cp-enterprise-kafka:5.2.3

python 客户端:

from kafka import KafkaProducer
import json
producer = KafkaProducer(bootstrap_servers=['localhost:29092'],
value_serializer=lambda v: json.dumps(v).encode('utf-8'),
security_protocol='PLAINTEXT')
acc_ini = 523416
print("Sending message")
producer.send('test', {'model_id': '1','acc':str(acc_ini), 'content':'test'})
producer.flush()

【讨论】:

  • host.docker.internal 不正确,因为您应该使用 --network 并连接到 zookeeper:2181(这是使用 Compose 文件所做的事情)。另外,Confluent 镜像不使用KAFKA_CREATE_TOPICS env-var,你的内部监听器应该是PLAINTEXT 才能在 docker 网络中准确配置
  • 另外,这并不能回答有关 Docker 容器中 Python 客户端的实际问题
  • 感谢您的 cmets @OneCricketeer。当然,您的建议也有效,但我的建议也有效,至少就我的目的而言,它将外部 python 消费者与 Kafka 容器连接起来。我解释说,在发布代码之前的第一件事......因为我在这个线程结束寻找我的问题的解决方案(外部 python 消费者),我认为可能还有其他人遇到我同样的问题也在这里结束。如果您发现任何其他线程更适合此评论,请告诉我。我没找到。
  • 我经常关闭这个stackoverflow.com/questions/51630260/… 的帖子,客户端语言无关紧要。也有很多博客。 confluent.io/blog/…
【解决方案4】:

在我的本地设置中,我遇到了同样的问题,所有容器在 docker 内都可以正常工作,但连接器无法连接到 Kafka。

我在docker-compose.yml里面的Kafka配置:

  broker:
    image: confluentinc/cp-server:7.0.1
    hostname: broker
    container_name: broker
    depends_on:
      - zookeeper
    ports:
      - "29092:29092"
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_LISTENERS: INTERNAL://:29092,EXTERNAL://:9092
      KAFKA_ADVERTISED_LISTENERS: INTERNAL://:29092,EXTERNAL://:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL

根据 google 的回复,我尝试将 localhost:port127.0.0.1:port0.0.0.0:port 作为连接器配置中的 URL,但没有任何效果。

最终,我必须传递本地系统的实际 IP 地址,并且它起作用了。 我在 Windows 10 上使用 Docker。

ActiveMQ 连接器配置文件:

{
  "connector.class": "io.confluent.connect.activemq.ActiveMQSourceConnector",
  "activemq.url": "tcp://<my-system-ip>:61616", <-- my system's IP address
  "max.poll.duration": "60000",
  "tasks.max": "1",
  "batch.size": "1",
  "name": "activemq-jms-connector",
  "jms.destination.name": "jms-test",
  "kafka.topic": "topic-1",
  "activemq.password": "password",
  "jms.destination.type": "topic",
  "use.permissive.schema": "false",
  "activemq.username": "username"
}

我希望这将帮助任何在 Windows 和 Docker 上遇到 Kafka-connect 设置的人。

【讨论】:

    猜你喜欢
    • 2018-04-28
    • 1970-01-01
    • 2018-08-04
    • 2019-03-24
    • 2014-07-30
    • 1970-01-01
    • 2018-04-02
    • 1970-01-01
    • 2019-07-01
    相关资源
    最近更新 更多