Debezium SQL Server连接器schema_only模式下未知记录问题问询
Debezium SQL Server连接器配置问题
我在配置Debezium SQL Server连接器时,将snapshot.mode设为schema_only,发现即使事件日志里没有新条目,Debezium在流启动后也会立刻发送一条记录。这条记录导致ExtractField$Key转换报错,因为它的key里没有目标字段payment_id。我通过首次运行移除转换、之后再添加的方式规避了问题,但即使移除转换,这条记录也没发到业务topic,没法查看内容,而正常记录的key提取是正常的。
我有三个问题:
- 这条记录是什么,为什么会被发送?
- 移除转换后,该记录为何没有到达业务topic?
- 是否有办法记录Debezium的所有输出记录?
连接器日志
[2022-11-16 08:41:16,292] INFO [inquiry_events_testing_8|task-0] No previous offset has been found (io.debezium.connector.sqlserver.SqlServerSnapshotChangeEventSource:68) [2022-11-16 08:41:16,292] INFO [inquiry_events_testing_8|task-0] According to the connector configuration only schema will be snapshotted (io.debezium.connector.sqlserver.SqlServerSnapshotChangeEventSource:73) [2022-11-16 08:41:16,292] INFO [inquiry_events_testing_8|task-0] Snapshot step 1 - Preparing (io.debezium.relational.RelationalSnapshotChangeEventSource:90) [2022-11-16 08:41:16,499] INFO [inquiry_events_testing_8|task-0] Snapshot step 2 - Determining captured tables (io.debezium.relational.RelationalSnapshotChangeEventSource:99) [2022-11-16 08:41:16,591] INFO [inquiry_events_testing_8|task-0] Adding table PCH.dbo.DATABASECHANGELOG to the list of capture schema tables (io.debezium.relational.RelationalSnapshotChangeEventSource:192) [2022-11-16 08:41:16,592] INFO [inquiry_events_testing_8|task-0] Adding table PCH.dbo.DATABASECHANGELOGLOCK to the list of capture schema tables (io.debezium.relational.RelationalSnapshotChangeEventSource:192) [2022-11-16 08:41:16,592] INFO [inquiry_events_testing_8|task-0] Adding table PCH.dbo.event_log to the list of capture schema tables (io.debezium.relational.RelationalSnapshotChangeEventSource:192) [2022-11-16 08:41:16,592] INFO [inquiry_events_testing_8|task-0] Snapshot step 3 - Locking captured tables [PCH.dbo.event_log] (io.debezium.relational.RelationalSnapshotChangeEventSource:106) [2022-11-16 08:41:16,592] INFO [inquiry_events_testing_8|task-0] Setting locking timeout to 10 s (io.debezium.connector.sqlserver.SqlServerSnapshotChangeEventSource:126) [2022-11-16 08:41:16,596] INFO [inquiry_events_testing_8|task-0] Executing schema locking (io.debezium.connector.sqlserver.SqlServerSnapshotChangeEventSource:131) [2022-11-16 08:41:16,596] INFO [inquiry_events_testing_8|task-0] Locking table PCH.dbo.event_log (io.debezium.connector.sqlserver.SqlServerSnapshotChangeEventSource:138) [2022-11-16 08:41:16,598] INFO [inquiry_events_testing_8|task-0] Snapshot step 4 - Determining snapshot offset (io.debezium.relational.RelationalSnapshotChangeEventSource:112) [2022-11-16 08:41:16,600] INFO [inquiry_events_testing_8|task-0] Snapshot step 5 - Reading structure of captured tables (io.debezium.relational.RelationalSnapshotChangeEventSource:115) [2022-11-16 08:41:16,693] INFO [inquiry_events_testing_8|task-0] Reading structure of schema 'PCH' (io.debezium.connector.sqlserver.SqlServerSnapshotChangeEventSource:199) [2022-11-16 08:41:16,796] INFO [inquiry_events_testing_8|task-0] Snapshot step 6 - Persisting schema history (io.debezium.relational.RelationalSnapshotChangeEventSource:119) [2022-11-16 08:41:16,796] WARN [inquiry_events_testing_8|task-0] The Kafka Connect schema name 'oi-pc-cdc-payment-case-handling-inquiry-events-1.dbo.event_log.Value' is not a valid Avro schema name, so replacing with 'oi_pc_cdc_payment_case_handling_inquiry_events_1.dbo.event_log.Value' (io.debezium.util.SchemaNameAdjuster:172) [2022-11-16 08:41:16,796] WARN [inquiry_events_testing_8|task-0] The Kafka Connect schema name 'oi-pc-cdc-payment-case-handling-inquiry-events-1.dbo.event_log.Key' is not a valid Avro schema name, so replacing with 'oi_pc_cdc_payment_case_handling_inquiry_events_1.dbo.event_log.Key' (io.debezium.util.SchemaNameAdjuster:172) [2022-11-16 08:41:16,797] WARN [inquiry_events_testing_8|task-0] The Kafka Connect schema name 'oi-pc-cdc-payment-case-handling-inquiry-events-1.dbo.event_log.Envelope' is not a valid Avro schema name, so replacing with 'oi_pc_cdc_payment_case_handling_inquiry_events_1.dbo.event_log.Envelope' (io.debezium.util.SchemaNameAdjuster:172) [2022-11-16 08:41:17,098] INFO [inquiry_events_testing_8|task-0] Schema locks released. (io.debezium.connector.sqlserver.SqlServerSnapshotChangeEventSource:158) [2022-11-16 08:41:17,098] INFO [inquiry_events_testing_8|task-0] Snapshot step 7 - Skipping snapshotting of data (io.debezium.relational.RelationalSnapshotChangeEventSource:135) [2022-11-16 08:41:17,100] INFO [inquiry_events_testing_8|task-0] Snapshot - Final stage (io.debezium.pipeline.source.AbstractSnapshotChangeEventSource:88) [2022-11-16 08:41:17,101] INFO [inquiry_events_testing_8|task-0] Removing locking timeout (io.debezium.connector.sqlserver.SqlServerSnapshotChangeEventSource:239) [2022-11-16 08:41:17,103] INFO [inquiry_events_testing_8|task-0] Snapshot ended with SnapshotResult [status=COMPLETED, offset=SqlServerOffsetContext [sourceInfoSchema=Schema{io.debezium.connector.sqlserver.Source:STRUCT}, sourceInfo=SourceInfo [serverName=oi-pc-cdc-payment-case-handling-inquiry-events-1, changeLsn=NULL, commitLsn=000002a6:00005ba0:0001, eventSerialNo=null, snapshot=FALSE, sourceTime=2022-11-16T08:41:16.796Z], snapshotCompleted=true, eventSerialNo=1]] (io.debezium.pipeline.ChangeEventSourceCoordinator:156) [2022-11-16 08:41:17,103] INFO [inquiry_events_testing_8|task-0] Connected metrics set to 'true' (io.debezium.pipeline.ChangeEventSourceCoordinator:236) [2022-11-16 08:41:17,103] INFO [inquiry_events_testing_8|task-0] Starting streaming (io.debezium.connector.sqlserver.SqlServerChangeEventSourceCoordinator:87) [2022-11-16 08:41:17,103] INFO [inquiry_events_testing_8|task-0] Last position recorded in offsets is 000002a6:00005ba0:0001(NULL)[1] (io.debezium.connector.sqlserver.SqlServerStreamingChangeEventSource:139) [2022-11-16 08:41:17,202] ERROR [inquiry_events_testing_8|task-0] WorkerSourceTask{id=inquiry_events_testing_8-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask:195) org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:223) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:149) at org.apache.kafka.connect.runtime.TransformationChain.apply(TransformationChain.java:50) at org.apache.kafka.connect.runtime.WorkerSourceTask.sendRecords(WorkerSourceTask.java:355) at org.apache.kafka.connect.runtime.WorkerSourceTask.execute(WorkerSourceTask.java:258) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:188) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:243) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: java.lang.IllegalArgumentException: Unknown field: payment_id at org.apache.kafka.connect.transforms.ExtractField.apply(ExtractField.java:65) at org.apache.kafka.connect.runtime.TransformationChain.lambda$apply$0(TransformationChain.java:50) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:173) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:207) ... 11 more [2022-11-16 08:41:17,203] INFO [inquiry_events_testing_8|task-0] Stopping down connector (io.debezium.connector.common.BaseSourceTask:243)
连接器配置
{ "name": "inquiry_events_testing_8", "config": { "connector.class": "io.debezium.connector.sqlserver.SqlServerConnector", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "key.converter.schemas.enable": "false", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "database.hostname": "${db_cluster}", "database.dbname" : "PCH", "database.user": "${db_user}", "database.password": "${db_password}", "table.include.list": "dbo.event_log", "database.history.kafka.bootstrap.servers": "kafka.core.svc:9092", "database.history.kafka.topic": "oi-pc-cdc-payment-case-handling-inquiry-events-db-history", "database.server.name": "${db_server_name}", "snapshot.mode": "schema_only", "message.key.columns": "dbo.event_log:payment_id", "transforms": "Reroute, Unwrap, ExtractValueFrom", "transforms.Reroute.type": "io.debezium.transforms.ByLogicalTableRouter", "transforms.Reroute.topic.regex": "${db_server_name}.dbo.event_log", "transforms.Reroute.topic.replacement": "${sink_topic}", "transforms.Unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.ExtractValueFrom.type": "org.apache.kafka.connect.transforms.ExtractField$Key", "transforms.ExtractValueFrom.field": "payment_id" } }
问题解答
1. 这条记录是什么,为什么会被发送?
这条记录是Debezium在schema_only快照完成后发送的元数据事件(大概率是表结构同步事件或快照结束标记)。当使用snapshot.mode=schema_only时,Debezium会在完成表结构读取并持久化到schema history后,发送一条事件来同步表结构信息,或者通知系统快照阶段结束。这类事件不属于业务数据变更,不会遵循message.key.columns配置的key规则,它的key可能仅包含服务器、表名等元信息,甚至是空结构,所以ExtractField$Key转换找不到payment_id字段会报错。
2. 移除转换后,该记录为何没有到达业务topic?
有两个核心原因:
- 这类元数据事件默认会发送到你配置的database history topic(
oi-pc-cdc-payment-case-handling-inquiry-events-db-history),而非业务sink topic,你之前查看的是业务topic,所以看不到。 - 若移除转换前连接器因报错崩溃,这条事件可能未完成Kafka事务提交,消费者无法读取;移除转换后连接器正常启动,但该事件仅在首次快照完成时发送一次,后续不会重复发送,你后续查看topic自然找不到。
3. 是否有办法记录Debezium的所有输出记录?
可以通过以下方式捕获所有Debezium输出的记录:
- 调高声量日志级别:将
io.debezium和org.apache.kafka.connect的日志级别设为DEBUG,连接器生成的所有事件(包括元数据事件、业务事件)都会被打印到日志文件中,直接查看日志就能找到目标记录的结构。 - 临时使用FileSink:配置一个FileStreamSink连接器,将所有输出路由到本地文件,这样所有事件都会被写入文件,包括那条触发报错的元数据事件。
- 消费database history topic:直接消费你配置的schema历史topic,元数据事件通常会发送到这里,查看该topic的消息即可获取事件内容。
内容的提问来源于stack exchange,提问作者Mindaugas Volskas
相关产品推荐
相关产品推荐

