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

Oracle数据库下Kafka Connect的SMT ExtractField提取字段报错咨询

Oracle Debezium连接器ExtractField SMT提取字段报错问题

问题描述

我在使用Oracle数据库时,无法执行SMT转换ExtractField将键结构体中的字段提取为简单long值,相同操作在Postgres数据库中可正常运行。
我尝试使用ReplaceField SMT重命名键,功能运行正常。我怀疑org.apache.kafka.connect.transforms.ExtractField类在获取字段的Schema处理逻辑存在问题,ReplaceField和ExtractField的Schema处理机制似乎存在差异。

相关版本信息

  • Oracle数据库版本:Oracle Database 19c Enterprise Edition Release 19.0.0.0.0 - Production Version 19.8.0.0.0
  • Debezium connect版本:1.6
  • Kafka版本:2.7.0
  • Instanclient basic(Oracle客户端及驱动)版本:21.3.0.0.0

错误信息

运行时抛出Unknown field ID_MYTABLE错误,错误栈如下:

org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:206) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:132) at org.apache.kafka.connect.runtime.TransformationChain.apply(TransformationChain.java:50) at org.apache.kafka.connect.runtime.WorkerSourceTask.sendRecords(WorkerSourceTask.java:339) at org.apache.kafka.connect.runtime.WorkerSourceTask.execute(WorkerSourceTask.java:264) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:185) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:234) 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:834) Caused by: java.lang.IllegalArgumentException: Unknown field: ID_MYTABLE 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:156) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:190) ... 11 more

连接器配置

{
  "name": "oracle-connector",
  "config": {
    "connector.class": "io.debezium.connector.oracle.OracleConnector",
    "tasks.max": "1",
    "database.server.name": "serverName",
    "database.user": "c##dbzuser",
    "database.password": "dbz",
    "database.url": "jdbc:oracle:thin:...",
    "database.dbname": "dbName",
    "database.pdb.name": "PDBName",
    "database.connection.adapter": "logminer",
    "database.history.kafka.bootstrap.servers": "kafka:9092",
    "database.history.kafka.topic": "schema-changes.data",
    "schema.include.list": "mySchema",
    "table.include.list": "mySchema.myTable",
    "log.mining.strategy": "online_catalog",
    "snapshot.mode": "initial",
    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "key.converter.schemas.enable": "false",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schemas.enable": "true",
    "value.converter.schema.registry.url": "http://schema-registry:8081",
    "transforms": "unwrap,route,extractField",
    "transforms.extractField.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
    "transforms.extractField.field": "ID_MYTABLE",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
    "transforms.route.regex": "([^.]+)\\.([^.]+)\\.([^.]+)",
    "transforms.route.replacement": "$1_$2_$3"
  }
}

问题原因

核心原因是ExtractField和ReplaceField的实现逻辑差异:

  • ReplaceField不需要依赖Schema,可以直接操作消息负载的键值对,所以关闭Key Schema的情况下也能正常运行
  • ExtractField严格依赖Schema来定位字段,你当前配置中key.converter.schemas.enable设为false,Key的Schema信息被丢弃,导致ExtractField无法识别结构体中的字段,抛出未知字段错误。
    另外Debezium Oracle连接器默认会将表字段名转为大写,也可能存在实际主键字段名和配置的ID_MYTABLE不匹配的情况。

解决方案

方案1:开启Key Schema支持(推荐)

修改连接器配置,将key.converter.schemas.enable改为true,如果需要和值端格式统一,也可以将Key转换器替换为AvroConverter,配置如下:

"key.converter": "io.confluent.connect.avro.AvroConverter",
"key.converter.schema.registry.url": "http://schema-registry:8081",
"key.converter.schemas.enable": "true"

方案2:验证字段实际命名

可以临时添加调试SMT输出Key的实际结构,确认主键字段的真实名称,也可以添加Debezium配置database.field.name.adjustment调整字段大小写策略,避免大小写不匹配问题。

方案3:无Schema替代方案

如果不需要保留Key Schema,可以替换ExtractField为自定义SMT或者支持无Schema提取的转换组件,直接从消息负载中提取对应键值,不需要依赖Schema信息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 22:27:02