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

Debezium SQL Server连接器schema_only模式下未知记录问题问询

Debezium SQL Server连接器配置问题

我在配置Debezium SQL Server连接器时,将snapshot.mode设为schema_only,发现即使事件日志里没有新条目,Debezium在流启动后也会立刻发送一条记录。这条记录导致ExtractField$Key转换报错,因为它的key里没有目标字段payment_id。我通过首次运行移除转换、之后再添加的方式规避了问题,但即使移除转换,这条记录也没发到业务topic,没法查看内容,而正常记录的key提取是正常的。

我有三个问题:

  1. 这条记录是什么,为什么会被发送?
  2. 移除转换后,该记录为何没有到达业务topic?
  3. 是否有办法记录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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 20:20:28