【问题标题】:BrokerConnect Issue [duplicate]BrokerConnect 问题 [重复]
【发布时间】:2021-04-30 07:34:21
【问题描述】:

我有一个用 python 编写的 kafka 生产者,我已将其添加到 docker-compose.yml

制作人:

import os, csv, avro.schema, json
from avro.datafile import DataFileReader, DataFileWriter
from avro.io import DatumReader, DatumWriter
from kafka import KafkaProducer
from collections import namedtuple

output_loc = '{}/avro.avro'.format(os.path.dirname(__file__))
CSV = '{}/oscar_age_male.csv'.format(os.path.dirname(__file__))
fields = ("Index","Year", "Age", "Name", "Movie")
csv_record = namedtuple('csv_record', fields)

p = KafkaProducer(bootstrap_servers = ['localhost:9092', 'kafka:9092'], value_serializer=lambda v: json.dumps(v).encode('utf-8'))

def read_csv(path):
    with open(path, 'rU') as data:
        data.readline()
        reader = csv.reader(data, delimiter=",")
        for row in map(csv_record._make, reader):
            yield row

def parse_schema(path='{}/schema.avsc'.format(os.path.dirname(__file__))):
    with open(path, 'r') as data:
        return avro.schema.parse(data.read())

def serilialise_records_and_send(records, outpath=output_loc):
    schema = parse_schema()
    with open(outpath, 'w') as out:
        writer = DataFileWriter(out, DatumWriter(), schema)
        for record in records:
            record = dict((f, getattr(record, f)) for f in record._fields)
            writer.append(record)
            msg = p.send(topic='test', value=record)
            metadt = msg.get()
            print(metadt.topic)
            print(metadt.partition)

serilialise_records_and_send(read_csv(CSV))

当我运行 docker-compose 时,由于没有可用的代理,我的生产者映像失败。

谁能告诉我为什么经纪人不可用?

当我从 IDE 本地运行生产者时,我可以连接,所以不确定缺少什么

docker-compose.yml

version: '2'
services:
  zookeeper:
    image: "confluentinc/cp-zookeeper:5.4.0"
    hostname: zookeeper
    ports:
      - '32181:32181'
    environment:
      ZOOKEEPER_CLIENT_PORT: 32181
      ZOOKEEPER_TICK_TIME: 2000
    extra_hosts:
      - "moby:127.0.0.1"  
      
  kafka:
    image: "confluentinc/cp-enterprise-kafka:5.4.0"
    hostname: kafka
    ports:
      - '9092:9092'
      - '29092:29092'
    depends_on:
      - zookeeper
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:32181
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:29092,PLAINTEXT_HOST://0.0.0.0:9092
      KAFKA_METRIC_REPORTERS: io.confluent.metrics.reporter.ConfluentMetricsReporter
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      CONFLUENT_METRICS_REPORTER_BOOTSTRAP_SERVERS: kafka:29092
      CONFLUENT_METRICS_REPORTER_ZOOKEEPER_CONNECT: zookeeper:32181
      CONFLUENT_METRICS_REPORTER_TOPIC_REPLICAS: 1
      CONFLUENT_METRICS_ENABLE: 'false'
      CONFLUENT_SUPPORT_CUSTOMER_ID: 'anonymous'
    extra_hosts:
      - "moby:127.0.0.1"  
      
  schema-registry:
    image: "confluentinc/cp-schema-registry:latest"
    hostname: schema-registry
    depends_on:
      - zookeeper
      - kafka
    ports:
      - '8081:8081'
    environment:
      SCHEMA_REGISTRY_HOST_NAME: schema-registry
      SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL: zookeeper:32181
      SCHEMA_REGISTRY_ACCESS_CONTROL_ALLOW_METHODS: GET,POST,PUT,OPTIONS
      SCHEMA_REGISTRY_ACCESS_CONTROL_ALLOW_ORIGIN: '*' 
    extra_hosts:
      - "moby:127.0.0.1"

      
  kafdrop:
    image: "obsidiandynamics/kafdrop"
    ports:
      - '9000:9000'
    environment:
      KAFKA_BROKERCONNECT: kafka:29092
      JVM_OPTS: "-Xms32M -Xmx64M"
      SERVER_SERVLET_CONTEXTPATH: "/"
    depends_on:
      - kafka
  
  producer:
    image: "producer"
    ports: 
      - '5000:5000'
    environment: 
      KAFKA_BROKERCONNECT: kafka:29092
    depends_on: 
      - kafka

【问题讨论】:

    标签: python docker apache-kafka


    【解决方案1】:

    您已经为 Kafka 定义了一个环境变量。你应该使用它

    import os
    p = KafkaProducer(bootstrap_servers = [os.environ['KAFKA_BROKERCONNECT']])
    

    当代码在容器中运行时,您不能使用localhost,因为它指的是那个服务,而不是代理


    但是,对您显示的代码更通用的解决方案是使用 CSV SpoolDir Connector

    【讨论】:

      【解决方案2】:

      您正在硬编码代理主机和端口

      p = KafkaProducer(bootstrap_servers = ['localhost:9092', 'kafka:9092'], value_serializer=lambda v: json.dumps(v).encode('utf-8'))
      

      但是在 Docker 下运行时,您需要连接到适合该网络的侦听器,即kafka:29092(您已将其包含在KAFKA_BROKERCONNECT,但我没有看到它在您的代码中被读取)。

      所以要么更新代码以使用环境变量,要么将硬编码列表更改为

      p = KafkaProducer(bootstrap_servers = ['localhost:9092', 'kafka:29092'], value_serializer=lambda v: json.dumps(v).encode('utf-8'))
      

      请参阅this blog,了解有关 Kafka 侦听器、Docker 等的全面背景。

      【讨论】:

      • 进行了建议的更改,但问题仍然存在
      • 你读过博客吗?
      • 是的,我似乎无法解决这个问题我会继续研究谢谢。
      • 自行尝试bootstrap_servers = ['kafka:29092']
      猜你喜欢
      • 1970-01-01
      • 2013-02-02
      • 2021-11-02
      • 1970-01-01
      • 2011-08-26
      • 2011-10-13
      • 2011-05-29
      • 2021-04-17
      • 2018-05-19
      相关资源
      最近更新 更多