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

Docker容器启动时Consumer脚本未自动执行问题求助

Kafka消费者容器无日志输出问题

问题背景

刚学习Docker,正在搭建一个将生产者、消费者和服务器脚本分别部署在独立容器中的项目。消费者负责更新/app/data目录下sqlite3数据库中的表,希望通过API请求访问该数据库。

预期效果

Kafka、Zookeeper和Server容器完成健康检查后,生产者和消费者容器正常运行,且能通过主机的localhost:5000向服务器容器发送API请求。

实际情况

Kafka、Zookeeper、Server和生产者容器均正常工作,但消费者容器启动后无任何日志输出,也没有执行Dockerfile中CMD指令指定的脚本。

已尝试操作

  • 所有脚本在容器内单独运行正常;
  • 使用docker exec -it命令在消费者容器内执行Python脚本也能正常运行;
  • 修改端口配置后容器会输出错误日志;
  • 将consumer.py替换为Hello World脚本后,容器能正常执行该脚本。

因此怀疑问题出在consumer.py本身,尽管它在容器外配置正确端口时能正常运行。

相关代码

消费者Python脚本 (consumer.py)

from confluent_kafka import Consumer, KafkaException , KafkaError
import json
import sqlite3
import uuid

# Kafka consumer configuration

conf = {
    'bootstrap.servers': 'kafka:9092',
    'group.id': uuid.uuid4,
    'auto.offset.reset': 'earliest',
    'enable.auto.commit': False  # Disable auto commit to manually commit offsets
}

consumer = Consumer(conf)
consumer.subscribe(['pokemon'])

# SQLite database configuration
conn = sqlite3.connect('/app/data/pokemon.db')
cursor = conn.cursor()



def insert_into_db(name, description, price, stock):
    cursor.execute('''
    INSERT OR REPLACE INTO pokemon (name, description, price, stock) VALUES (?, ?, ?, ?)
    ''', (name, description, price, stock))
    conn.commit()


try:
    while True:
        msg = consumer.poll(timeout=1.0)
        if msg is None:
            continue
        if msg.error():
            if msg.error().code() == KafkaError._PARTITION_EOF:
                print('End of partition reached {0}/{1}'.format(msg.topic(), msg.partition()))
            elif msg.error():
                raise KafkaException(msg.error())
        else:
            message = json.loads(msg.value().decode('utf-8'))
            name = message.get('name')
            print("Received information of " + name)
            description = message.get('Description')
            price = message.get('price')
            stock = message.get('stock')

            # Insert message data into database
            insert_into_db(name, description, price, stock)
            print([name, description, price, stock])
            # Commit offset after writing to file
            consumer.commit(msg)
except Exception as e:
    print(f"An error occurred: {e}")
finally:
    consumer.close()
    conn.close()

print("Consumer closed and database connection closed.")

Docker Compose配置

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:latest
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
    ports:
      - "2181:2181"
    healthcheck:
      test: ["CMD", "nc", "-z", "localhost", "2181"]
      interval: 10s
      retries: 3
      start_period: 10s
      timeout: 5s
    
  kafka:
    image: confluentinc/cp-kafka:latest
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
      - "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
      


    healthcheck:
      test: ["CMD", "nc", "-z", "localhost", "9092"]
      interval: 10s
      retries: 3
      start_period: 10s
      timeout: 5s

  producer:
    build: ./producer
    depends_on:
      kafka:
        condition: service_healthy


  consumer:
    build: ./consumer
    depends_on:
      producer:
        condition: service_completed_successfully
      server:
         condition: service_healthy
    volumes:
      - db-data:/app/data  

  server:
    build: ./server

    depends_on:
      - kafka
      
    environment:
      KAFKA_BOOTSTRAP_SERVERS: kafka:9092  # Add environment variable for Kafka servers
    volumes:
      - db-data:/app/data  # Mount the shared volume for the database
    ports:
      - "5000:5000"  # Expose port for the Flask app
    healthcheck:
      test: ["CMD-SHELL", "[ -f /app/data/pokemon.db ]"]  # Check if the database file exists
      interval: 10s
      retries: 3
      start_period: 30s
      timeout: 5s

 
volumes:
  db-data:

消费者Dockerfile

FROM python:3.11


WORKDIR /app


RUN pip install confluent-kafka

COPY . /app/

CMD ["python", "./consume.py"]

内容的提问来源于stack exchange,提问作者barkin sagan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 01:00:55