Oracle数据库下Kafka Connect的SMT ExtractField提取字段报错咨询
问题描述
我在使用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

