You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求助: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 16:54:55