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

使用Debezium PostgreSQL Connector时Topic为空问题求助

问题排查:Debezium PostgreSQL Connector日志显示发送记录但Kafka Topic无数据

核心场景

Windows环境下部署Kafka + Debezium PostgreSQL Connector,向PostgreSQL的public.student表插入数据时,Connector日志显示已发送记录,但对应Kafka Topic中无数据。已将Debezium PostgreSQL的jar包放入Kafka/libs目录。

排查与解决步骤

1. 确认目标Topic名称是否正确

Debezium PostgreSQL Connector生成的Topic命名规则为:{database.server.name}.{schema}.{table}。
你的配置中database.server.name = postgres,监听表为public.student,因此正确的Topic名称应为postgres.public.student,而非你可能误以为的fulfillment相关名称(配置中同时设置的topic.prefix是旧版参数,与database.server.name冲突,会被后者覆盖)。

2. 验证Kafka Topic状态与消息

使用Kafka命令行工具检查:

  • 列出所有Topic,确认目标Topic存在:
    .\kafka\bin\windows\kafka-topics.bat --list --bootstrap-server localhost:9092
    
  • 消费目标Topic的全部消息,插入新数据后观察是否能收到:
    .\kafka\bin\windows\kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic postgres.public.student --from-beginning
    

3. 检查Kafka Connect配置正确性

打开connect-standalone.properties,确认以下配置:

  • bootstrap.servers是否指向正确的Kafka地址(默认应为localhost:9092),确保Connect能正常连接Kafka集群;
  • key.converter和value.converter配置正常(例如org.apache.kafka.connect.json.JsonConverter),无格式错误导致消息无法写入。

4. 查看Kafka Connect完整日志

除了已看到的发送记录日志,检查是否有报错信息:

  • 确认connect-standalone.properties中日志配置(如log4j.rootLogger)级别为INFO或DEBUG,以便查看更详细的消息发送过程;
  • 排查是否存在网络异常、Topic自动创建关闭(auto.create.topics.enable是否为true)等导致消息写入失败的情况。

5. 验证PostgreSQL配置与复制槽

  • 确认PostgreSQL的postgresql.conf中wal_level = logical,并已重启PostgreSQL生效(逻辑复制是Debezium工作的前提);
  • 登录PostgreSQL执行以下SQL,检查复制槽状态:
    SELECT slot_name, active FROM pg_replication_slots;
    
    确保debezium复制槽存在且active为t(表示Connector正在正常使用该槽)。

6. 清理冲突的Connector配置

你的配置中同时设置了database.server.name和topic.prefix,这两个参数互斥(topic.prefix是旧版参数,对应新版的database.server.name)。建议删除topic.prefix配置项,避免命名混淆,修改后的配置示例:

{
  "name": "my-postgres-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "tasks.max": "1",
    "database.hostname": "localhost",
    "database.port": "5432",
    "database.user": "postgres",
    "database.password": "1524",
    "database.dbname": "MartenDB",
    "database.server.name": "postgres",
    "table.include.list": "public.student",
    "plugin.name": "pgoutput",
    "slot.name": "debezium"
  }
}

修改后重启Connector,重新发送配置。


内容的提问来源于stack exchange,提问作者Sidhant Suvagiya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 07:26:05