【问题标题】:I am trying to run a flask app and kafka through docker compose..but unable to set things right我正在尝试通过 docker compose 运行烧瓶应用程序和 kafka ..但无法正确设置
【发布时间】:2020-02-06 13:35:39
【问题描述】:

我正在尝试通过 docker compose..along 和 kafka 运行烧瓶应用程序。我正在一起运行烧瓶应用程序和 consumer.py。当我通过烧瓶 api 调用 producer.py 时,它似乎正在发送数据,但消费者没有收到任何东西。由于我同时运行烧瓶应用程序和消费者文件..我也无法查看日志...可以提供一些帮助..

我的 docker-compose 文件

version: '3'
services:
  zookeeper:
    image: wurstmeister/zookeeper
    container_name: "zoo1"
    ports:
      - "2181:2181"
  kafka:
    image: wurstmeister/kafka
    container_name: "kafka1"
    ports:
     - "9092:9092"
    expose:
     - "9093"
    depends_on:
     - zookeeper
    environment:
      KAFKA_ADVERTISED_LISTENERS: INSIDE://kafka:9093,OUTSIDE://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INSIDE:PLAINTEXT,OUTSIDE:PLAINTEXT
      KAFKA_LISTENERS: INSIDE://0.0.0.0:9093,OUTSIDE://0.0.0.0:9092
      KAFKA_INTER_BROKER_LISTENER_NAME: INSIDE
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_CREATE_TOPICS: "test:1:1"
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true'
    volumes:
      - /var/run/docker.sock:/var/run/docker.sock
  backend:
    build: .
    image: backend:v1
    command: bash -c "python3 run.py run && python run.py consumer"
    container_name: backend
    restart: always
    volumes:
      - .:/app
    ports:
      - '5000:5000'
    depends_on:
      - mongodb
      - neodb
      - elasticsearch
      - kafka

我的生产者.py

from kafka import KafkaProducer
from json import dumps

class Producer:
    def __init__(self):
        pass

    def producer(self, topic,data):

        producer = KafkaProducer(
        value_serializer=lambda m: dumps(m).encode('utf-8'),
        bootstrap_servers=['kafka:9092','kafka:9093','172.17.0.1:32783','172.17.0.1:32782','172.17.0.1:32781'])
        producer.send(topic, data)

我的消费者.py

from kafka import KafkaConsumer
from json import loads
import os

class Consumer:
    def consumer(self,topic):
        print("hello")
        consumer = KafkaConsumer(
                topic,
                auto_offset_reset='latest',
                enable_auto_commit=True,
                group_id='my-group-1',
                value_deserializer=lambda m: loads(m.decode('utf-8')),
                bootstrap_servers=['kafka:9092','kafka:9093','172.17.0.1:32783','172.17.0.1:32782','172.17.0.1:32781'])
        while True:
            dirName = "bad"
            if not os.path.exists(dirName):
                os.mkdir(dirName)
            for m in consumer:
                pass

【问题讨论】:

    标签: python docker flask


    【解决方案1】:

    您似乎没有收到来自consumer = KafkaCo... 的任何对象

    这是您必须确保输入正确的文档。

    Apache Kafka Docs

    【讨论】:

    • 尝试在 Consumer 中创建构造函数,除非您是静态调用 Consumer。
    • 我能够设置 kafka...我尝试通过 kafka shell 从生产者向消费者发送消息...工作正常...但问题是当我通过 producer.py 发送消息时。 .我无法接收..
    • 寻找但只找到文档说把producer.close()放在最后
    猜你喜欢
    • 2017-05-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-01-24
    • 1970-01-01
    • 2020-10-24
    • 1970-01-01
    相关资源
    最近更新 更多