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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 09:45:31