使用Debezium PostgreSQL Connector时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

