求助:Kafka+FastAPI+Docker微服务模板消费者无法消费消息
简介
我目前正在尝试使用Kafka和FastAPI构建一个模板,以便快速开发微服务模式的程序。
目标愿景
打造一个实现简易微服务基础设施的设计模式仓库,示例仅展示不同服务间的消息传递方式,让用户无需花费大量时间进行环境配置,即可轻松集成自定义代码。
开发动机
我查找了大量资料,但未能找到通用的简易示例,多数示例高度定制化,不具备普适性。
技术栈
- Kafka
- FastAPI
- Docker
欢迎其他实现方案
我刚接触微服务架构,若您有其他建议,欢迎告知,我非常乐意探索更多设计思路。
当前问题
我当前的模板包含Zookeeper、Kafka、消费者和生产者服务,但遇到了消费者无法消费生产者生成的消息的问题。通过命令docker-compose exec kafka kafka-console-consumer.sh --bootstrap-server kafka:9092 --topic my-topic --from-beginning已确认生产者能成功发布消息,但消费者无任何响应。
我已多次重写消费者代码、修改端口与Docker Compose配置,但仍无法定位问题,恳请各位提供排查建议。
我的文件夹结构:

我的docker-compose文件:
version: '3' services: zookeeper: image: confluentinc/cp-zookeeper:latest environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 ports: - 2181:2181 - 2888:2888 - 3888:3888 kafka: image: confluentinc/cp-kafka:latest restart: "no" links: - zookeeper ports: - 9092:9092 environment: KAFKA_BROKER_ID: 1 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_LISTENERS: INTERNAL://:29092,EXTERNAL://:9092 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka:29092,EXTERNAL://localhost:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL producer: build: ./producer ports: - '8000:8000' environment: - KAFKA_BOOTSTRAP_SERVERS=kafka:29092 depends_on: - kafka consumer: build: ./consumer environment: - KAFKA_BOOTSTRAP_SERVERS=kafka:29092 - KAFKA_GROUP_ID=my-group depends_on: - kafka kafdrop: image: obsidiandynamics/kafdrop restart: "no" environment: - KAFKA_BOOTSTRAP_SERVERS=kafka:29092 ports: - 9000:9000 depends_on: - kafka
我的生产者Dockerfile:
FROM python:3.8-slim-buster COPY . /app WORKDIR /app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000"]
我的生产者requirements.txt:
fastapi uvicorn confluent-kafka
我的生产者main.py:
import json from confluent_kafka import Producer from fastapi import FastAPI app = FastAPI() producer_conf = { 'bootstrap.servers': 'kafka:9092', 'client.id': 'my-app' } producer = Producer(producer_conf) def produce(data: dict): try: data = json.dumps(data).encode('utf-8') producer.produce('my-topic', value=data) producer.flush() return {"status": "success", "message": data} except Exception as e: return {"status": "error", "message": str(e)}
我的消费者Dockerfile:
FROM python:3.8-slim-buster COPY . /app WORKDIR /app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt CMD [ "python", "main.py" ]
我的消费者requirements.txt:
confluent-kafka
我的消费者main.py:
from confluent_kafka import Consumer, KafkaError conf = { 'bootstrap.servers': 'kafka:9092', 'auto.offset.reset': 'earliest', 'enable.auto.commit': True, 'group.id': 'my-group', 'api.version.request': True, 'api.version.fallback.ms': 0 } def consume_messages(): consumer = Consumer(conf) consumer.subscribe(['my-topic']) try: while True: msg = consumer.poll(1.0) if msg is None: continue if msg.error(): if msg.error().code() == KafkaError._PARTITION_EOF: print(f'Reached end of partition: {msg.topic()}[{msg.partition()}]') else: print(f'Error while consuming messages: {msg.error()}') else: print(f"Received message: {msg.value().decode('utf-8')}") except Exception as e: print(f"Exception occurred while consuming messages: {e}") finally: consumer.close() def startup(): consume_messages() if __name__ == "__main__": try: print("Starting consumer...") startup() except Exception as e: print(f"Exception occurred: {e}")
构建命令:
docker-compose up
触发生产者的curl命令:
curl -X POST http://localhost:8000/produce -H "Content-Type: application/json" -d '{"key": "nice nice nice"}'
我已多次重写消费者代码、修改端口与Docker Compose配置,但仍无法定位问题,恳请各位提供排查建议。
内容的提问来源于stack exchange,提问作者mm117
相关产品推荐
相关产品推荐

