使用kafka-python结合Debezium无法从Kafka消费消息的问题
问题:本地Python无法消费Kafka CDC消息,但Docker容器内正常
我搭建了完整的CDC数据管道,向PostgreSQL的student表插入示例数据后,在Docker容器内可以正常消费Kafka主题postgres.public.student的CDC消息,但本地运行Python代码时无法获取任何消息。
Docker Compose配置
version: "3.7" services: postgres: image: debezium/postgres:13 ports: - 5432:5432 environment: - POSTGRES_USER=docker - POSTGRES_PASSWORD=docker - POSTGRES_DB=exampledb zookeeper: image: confluentinc/cp-zookeeper:5.5.3 environment: ZOOKEEPER_CLIENT_PORT: 2181 kafka: image: confluentinc/cp-enterprise-kafka:5.5.3 depends_on: [zookeeper] environment: KAFKA_ZOOKEEPER_CONNECT: "zookeeper:2181" KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092 KAFKA_BROKER_ID: 1 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_JMX_PORT: 9991 ports: - 9092:9092 - 29092:29092 debezium: image: debezium/connect:1.4 environment: BOOTSTRAP_SERVERS: kafka:9092 GROUP_ID: 1 CONFIG_STORAGE_TOPIC: connect_configs OFFSET_STORAGE_TOPIC: connect_offsets KEY_CONVERTER: io.confluent.connect.avro.AvroConverter VALUE_CONVERTER: io.confluent.connect.avro.AvroConverter CONNECT_KEY_CONVERTER_SCHEMA_REGISTRY_URL: http://schema-registry:8081 CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL: http://schema-registry:8081 depends_on: [kafka] ports: - 8083:8083 schema-registry: image: confluentinc/cp-schema-registry:5.5.3 environment: - SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL: zookeeper:2181 - SCHEMA_REGISTRY_HOST_NAME: schema-registry - SCHEMA_REGISTRY_LISTENERS: http://schema-registry:8081,http://localhost:8081 ports: - 8081:8081 depends_on: [zookeeper, kafka]
Debezium连接器配置
{ "name": "exampledb-connector", "config": { "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "database.user": "docker", "database.dbname": "exampledb", "database.hostname": "postgres", "database.password": "docker", "name": "exampledb-connector", "database.server.name": "postgres", "table.include.list": "public.student", "plugin.name": "pgoutput", "database.port": "5432" }, "tasks": [ { "connector": "exampledb-connector", "task": 0 } ], "type": "source" }
验证命令
CURL命令
curl --location --request GET 'http://localhost:8083/connectors/exampledb-connector'
SQL命令
CREATE TABLE student ( id integer primary key, name varchar ) ALTER TABLE public.student REPLICA IDENTITY FULL INSERT INTO public.student (id, name) VALUES (1,'test2') INSERT INTO public.student (id, name) VALUES (2,'test2')
Docker容器内测试
通过以下命令可成功消费主题postgres.public.student的消息:
docker network ls docker run --tty --network debezium_default confluentinc/cp-kafkacat kafkacat -b kafka:9092 -C -s key=s -s value=avro -r http://schema-registry:8081 -t postgres.public.student
本地Python代码(无法获取消息)
try: import kafka import json import requests import os import sys from json import dumps from kafka import KafkaProducer from kafka import KafkaConsumer from confluent_kafka.schema_registry import SchemaRegistryClient from kafka import KafkaConsumer import json import requests import os import sys except Exception as e: pass SCHEME_REGISTERY = "http://schema-registry:8081" TOPIC = "postgres.public.student" BROKER = "localhost:9092" import kafka consumer = kafka.KafkaConsumer(group_id='1', bootstrap_servers=[BROKER]) print(consumer.topics()) print("***************") def main(): print("Listening *****************") consumer = KafkaConsumer( TOPIC, bootstrap_servers=[BROKER], auto_offset_reset='earliest', enable_auto_commit=False, group_id='1' ) for msg in consumer: payload = json.loads(msg.value) payload["meta_data"]={ "topic":msg.topic, "partition":msg.partition, "offset":msg.offset, "timestamp":msg.timestamp, "timestamp_type":msg.timestamp_type, "key":msg.key, } print(payload, end="\n") main()
问题原因与解决方案
1. Kafka Broker地址错误
Docker Compose中Kafka配置了两个监听地址:
PLAINTEXT://kafka:9092:容器内部通信使用PLAINTEXT_HOST://localhost:29092:本地主机连接使用
你的Python代码中用了localhost:9092,这是容器内部端口,本地无法正确连接,需要改为localhost:29092。
2. 未处理Avro格式消息
Debezium连接器使用Avro转换器,Kafka中的消息是Avro二进制格式,不是JSON。直接用json.loads(msg.value)会失败,必须用Schema Registry客户端反序列化。
3. Schema Registry地址错误
本地访问Schema Registry需要用http://localhost:8081,而不是容器内部域名http://schema-registry:8081。
4. 消费者组偏移量冲突
容器内测试可能已经用group_id='1'消费过消息,导致偏移量已到最新位置。本地代码换一个新的group_id(比如python-cdc-consumer),才能重新拉取历史消息。
修正后的Python代码
from confluent_kafka import Consumer, KafkaError from confluent_kafka.schema_registry import SchemaRegistryClient from confluent_kafka.schema_registry.avro import AvroDeserializer # 正确的配置 SCHEMA_REGISTRY_URL = "http://localhost:8081" TOPIC = "postgres.public.student" BROKER = "localhost:29092" GROUP_ID = "python-cdc-consumer" # 初始化Schema Registry客户端 schema_registry_conf = {'url': SCHEMA_REGISTRY_URL} schema_registry_client = SchemaRegistryClient(schema_registry_conf) # 初始化Avro反序列化器 avro_deserializer = AvroDeserializer(schema_registry_client=schema_registry_client) # 配置Kafka消费者 consumer_conf = { 'bootstrap.servers': BROKER, 'group.id': GROUP_ID, 'auto.offset.reset': 'earliest', 'enable.auto.commit': False } consumer = Consumer(consumer_conf) consumer.subscribe([TOPIC]) def main(): print("Listening for CDC messages...") while True: msg = consumer.poll(1.0) if msg is None: continue if msg.error(): if msg.error().code() == KafkaError._PARTITION_EOF: print(f"Partition {msg.partition()} reached end") else: print(f"Consumer error: {msg.error()}") continue # 反序列化Avro消息 try: key = avro_deserializer(msg.key(), is_key=True) value = avro_deserializer(msg.value(), is_key=False) payload = { "key": key, "value": value, "meta_data": { "topic": msg.topic(), "partition": msg.partition(), "offset": msg.offset(), "timestamp": msg.timestamp(), "timestamp_type": msg.timestamp_type() } } print(payload) except Exception as e: print(f"Deserialization failed: {e}") # 手动提交偏移量 consumer.commit(msg) if __name__ == "__main__": try: main() except KeyboardInterrupt: print("Stopping consumer...") finally: consumer.close()
额外注意事项
- 安装依赖:
pip install confluent-kafka[avro] - 确认本地端口29092、8081未被占用
- 若仍无法消费,可先删除消费者组:
kafka-consumer-groups --bootstrap-server localhost:29092 --delete --group python-cdc-consumer
内容的提问来源于stack exchange,提问作者Soumil Nitin Shah
相关产品推荐
相关产品推荐

