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

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:
    wal_level = logical
    max_wal_senders = 10
    max_replication_slots = 10
    
    重启PostgreSQL容器:
    docker 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 04:17:09