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
相关产品推荐
相关产品推荐

