Debezium无法从PostgreSQL向Kafka推送数据问题求助
Debezium PostgreSQL CDC 连接器卡在WAL恢复位置搜索阶段
已部署的Docker容器
启动了以下Debezium及相关组件的Docker镜像:
# ZOOKEEPER docker run -it --rm --name zookeeper -p 2181:2181 -p 2888:2888 -p 3888:3888 quay.io/debezium/zookeeper:2.2 # KAFKA docker run -it --rm --name kafka -p 9092:9092 --link zookeeper:zookeeper quay.io/debezium/kafka:2.2 # KAFKA-CONNECT docker run -it --rm --name connect -p 8083:8083 -e GROUP_ID=1 -e CONFIG_STORAGE_TOPIC=my_connect_configs -e OFFSET_STORAGE_TOPIC=my_connect_offsets -e STATUS_STORAGE_TOPIC=my_connect_statuses --link kafka:kafka quay.io/debezium/connect:2.2 # KAFKA-UI docker run -it -p 8080:8080 -e DYNAMIC_CONFIG_ENABLED=true provectuslabs/kafka-ui # POSTGRESQL docker run --name postgresql -p 5432:5432 -e POSTGRES_PASSWORD=password123 -d debezium/postgres
PostgreSQL连接器配置
通过以下curl命令创建连接器:
curl -i -X POST -H "Accept:application/json" \ -H "Content-Type:application/json" localhost:8083/connectors/ \ -d '{ "name": "psql-connector", "config": { "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "plugin.name": "pgoutput", "tasks.max": "1", "database.hostname": "192.168.29.219", "database.port": "5432", "database.user": "postgres", "database.dbname": "demo", "database.password": "postgres", "database.include.list": "demo", "table.include.list": "public.company", "include.schema.changes": "true", "topic.prefix": "psqlserver", "schema.history.internal.kafka.bootstrap.servers": "kafka:9092", "schema.history.internal.kafka.topic": "schemahistory.demo" } }'
问题现象
已通过http://localhost:8083/connectors确认连接器创建成功,但Kafka Connect进程卡在以下日志阶段:
Postgres|psqlserver|streaming Requested thread factory for connector PostgresConnector, id = psqlserver named = keep-alive [io.debezium.util.Threads] Postgres|psqlserver|streaming Creating thread debezium-postgresconnector-psqlserver-keep-alive [io.debezium.util.Threads] Postgres|psqlserver|streaming Searching for WAL resume position [io.debezium.connector.postgresql.PostgresStreamingChangeEventSource]
解决方案
1. 修正数据库密码配置
启动PostgreSQL容器时设置的密码是password123,但连接器配置中database.password为postgres,两者不一致,这是核心问题之一。修改连接器配置中的密码为password123,重新更新连接器:
curl -i -X PUT -H "Accept:application/json" \ -H "Content-Type:application/json" localhost:8083/connectors/psql-connector/config \ -d '{ "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "plugin.name": "pgoutput", "tasks.max": "1", "database.hostname": "192.168.29.219", "database.port": "5432", "database.user": "postgres", "database.dbname": "demo", "database.password": "password123", "database.include.list": "demo", "table.include.list": "public.company", "include.schema.changes": "true", "topic.prefix": "psqlserver", "schema.history.internal.kafka.bootstrap.servers": "kafka:9092", "schema.history.internal.kafka.topic": "schemahistory.demo" }'
2. 验证PostgreSQL逻辑复制配置
Debezium依赖PostgreSQL的逻辑复制功能,需确保以下配置正确:
- 进入PostgreSQL容器:
docker exec -it postgresql bash - 检查wal_level:
若结果不是psql -U postgres -c "show wal_level;"logical,修改/var/lib/postgresql/data/postgresql.conf:
重启PostgreSQL容器:wal_level = logical max_wal_senders = 10 max_replication_slots = 10docker restart postgresql
3. 确认数据库用户权限
确保postgres用户拥有REPLICATION权限:
psql -U postgres -c "ALTER USER postgres WITH REPLICATION;"
4. 检查网络连通性
如果Kafka Connect容器无法通过宿主机IP访问PostgreSQL,可改用容器名postgresql作为database.hostname(因容器在同一Docker网络,通过--link可直接解析),修改连接器配置中的database.hostname为postgresql。
5. 检查复制槽状态
在PostgreSQL中查看是否生成了对应的复制槽:
psql -U postgres -c "SELECT slot_name, plugin FROM pg_replication_slots;"
若存在以psqlserver为前缀的槽,说明逻辑复制通道已建立;若不存在,需排查权限或连接配置问题。
内容的提问来源于stack exchange,提问作者Sai Ram
相关产品推荐
相关产品推荐

