Debezium+PostgreSQL CDC项目notifying方法未执行问题排查
问题排查与解决方法
一、核心配置与环境验证
1. 数据库权限校验
PostgreSQL CDC用户必须拥有逻辑复制权限及监控表的查询权限,执行以下SQL补全权限:
-- 授予数据库复制权限 GRANT REPLICATION SLAVE ON DATABASE your_db_name TO your_user; -- 授予Schema使用权限(PostgreSQL 10+必填) GRANT USAGE ON SCHEMA your_schema TO your_user; -- 授予目标表查询权限 GRANT SELECT ON ALL TABLES IN SCHEMA your_schema TO your_user;
2. PostgreSQL WAL配置检查
修改postgresql.conf后重启数据库,确保以下参数配置正确:
wal_level = logical(必须为logical,默认replica不支持逻辑复制)max_wal_senders = 10(至少大于1,确保有可用的WAL发送进程)max_replication_slots = 5(数值需大于你创建的复制槽数量)
3. 复制槽状态确认
登录PostgreSQL执行查询,验证复制槽是否正常激活:
SELECT slot_name, plugin, active, restart_lsn FROM pg_replication_slots;
active字段必须为t,否则说明Debezium未成功连接到该槽restart_lsn若长期不变,说明无变更数据被复制
二、Debezium Engine代码与配置检查
1. 关键配置项核对
确保Connector配置包含以下必填参数(示例):
Configuration config = Configuration.create() .with("connector.class", "io.debezium.connector.postgresql.PostgresConnector") .with("offset.storage", "org.apache.kafka.connect.storage.MemoryOffsetBackingStore") .with("database.hostname", "localhost") .with("database.port", "5432") .with("database.user", "cdc_user") .with("database.password", "cdc_pass") .with("database.dbname", "test_db") .with("server.name", "test_server") .with("database.include.list", "public.test_table") // 必须指定监控表(schema.table格式) .with("slot.name", "debezium_slot") // 与自动创建的槽名完全一致 .with("snapshot.mode", "initial") // 根据业务需求设置快照模式 .build();
database.include.list不能写错,否则不会监听目标表slot.name必须和自动创建的复制槽名称完全匹配
2. Engine启动逻辑校验
确保Debezium Engine在独立线程中运行,且主线程不提前退出:
// 创建Engine实例 DebeziumEngine<ChangeEvent<JsonNode, JsonNode>> engine = DebeziumEngine.create(Json.class) .using(config) .notifying(record -> { // 回调逻辑需添加异常捕获,避免静默失败 try { System.out.println("Received change event: " + record.value()); } catch (Exception e) { e.printStackTrace(); } }) .build(); // 用线程池托管Engine,避免主线程阻塞或退出 ExecutorService executor = Executors.newSingleThreadExecutor(); executor.submit(engine); // 注册关闭钩子,确保优雅停止 Runtime.getRuntime().addShutdownHook(new Thread(() -> { try { engine.close(); executor.shutdown(); } catch (Exception e) { e.printStackTrace(); } }));
- 若直接在主线程调用
engine.run(),会阻塞主线程;若主线程提前退出,Engine也会终止 - 回调方法内部的未捕获异常会导致静默失败,必须添加try-catch打印异常
3. 依赖版本兼容性检查
确保Maven依赖版本匹配,示例配置:
<dependencies> <dependency> <groupId>io.debezium</groupId> <artifactId>debezium-api</artifactId> <version>1.9.7.Final</version> </dependency> <dependency> <groupId>io.debezium</groupId> <artifactId>debezium-connector-postgresql</artifactId> <version>1.9.7.Final</version> <scope>runtime</scope> </dependency> <dependency> <groupId>org.postgresql</groupId> <artifactId>postgresql</artifactId> <version>42.5.4</version> </dependency> <dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> <version>2.13.4.2</version> </dependency> </dependencies>
- Debezium版本需与PostgreSQL版本兼容(如Debezium 1.9支持PostgreSQL 9.6-14)
- PostgreSQL驱动版本需对应数据库版本,避免兼容性问题
三、控制台日志关键点分析
从日志中查找以下关键信息,定位问题:
- 确认复制槽创建成功:
Successfully created replication slot 'debezium_slot' - 确认数据库连接正常:
Connected to PostgreSQL database 'test_db' with user 'cdc_user' - 确认已开始监听变更:
Streaming changes from LSN 0/1234567 - 若出现
Waiting for keepalive packet,说明连接正常但暂无变更,可插入测试数据验证 - 若日志无上述启动信息,说明Engine未正确初始化,检查配置是否存在拼写错误
四、其他潜在问题
- 目标表无主键:Debezium默认依赖主键追踪变更,若无主键,需添加配置
database.history.store.only.captured.tables.ddl=true,或明确指定table.include.list - 快照未完成:若配置
snapshot.mode=initial,首次启动会执行全量快照,需等待快照完成后才会监听增量变更 - SSL连接问题:若数据库要求SSL,需添加配置
database.ssl.mode=require
内容的提问来源于stack exchange,提问作者Jesús Alberto Carrillo García
相关产品推荐
相关产品推荐

