Debezium与Kafka集成Oracle时遇Unsupported source data type: STRUCT错误
集成Debezium、Kafka与Oracle时出现以下错误:
org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception. Caused by: org.apache.kafka.connect.errors.ConnectException: Unsupported source data type: STRUCT
环境配置
- Kafka独立模式配置:
bootstrap.servers=localhost:9092 value.converter=org.apache.kafka.connect.json.JsonConverter key.converter=org.apache.kafka.connect.storage.StringConverter offset.storage.file.filename=/tmp/connect.offsets plugin.path=kafka_2.12-3.2.3/libs
- Oracle源连接器配置:
name=connector-test180 connector.class=io.debezium.connector.oracle.OracleConnector tasks.max=1 database.server.name=server1 database.hostname=***.**.0.1 database.port=1521 database.user=username database.password=password database.dbname=ORCLCDB database.pdb.name=ORCLPDB1 database.connection.adapter=logminer database.history.kafka.bootstrap.servers=localhost:9092 database.history.kafka.topic=schema-changes.inventory table.include.list=DEBEZIUM.CUSTOMER column.include.list=DEBEZIUM.CUSTOMER.ID,DEBEZIUM.CUSTOMER.CUSTOMER_ID,DEBEZIUM.CUSTOMER.STATUS,DEBEZIUM.CUSTOMER.FIRSTNAME,\ DEBEZIUM.CUSTOMER.MOBILENUMBER,DEBEZIUM.CUSTOMER.FATHERNAME,DEBEZIUM.CUSTOMER.MOTHERNAME,DEBEZIUM.CUSTOMER.CITY,DEBEZIUM.CUSTOMER.COUNTRY time.precision.mode=connect transforms=filter,route transforms.filter.type=io.debezium.transforms.Filter transforms.filter.language=jsr223.groovy transforms.filter.condition=value.source.table == 'CUSTOMER' transforms.filter.topic.regex=server1.DEBEZIUM.* transforms.route.type=org.apache.kafka.connect.transforms.RegexRouter transforms.route.regex=([^.]+)\\.([^.]+)\\.([^.]+) transforms.route.replacement=$3
- Kafka Connect Sink连接器配置:
name=customer-sink139 connector.class=io.confluent.connect.jdbc.JdbcSinkConnector tasks.max=1 topics=CUSTOMER connection.url=jdbc:mysql://localhost:3306/dbname connection.user=user connection.password=pass transforms=unwrap transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState transforms.unwrap.drop.tombstones=false auto.create=false insert.mode=upsert pk.mode=record_value errors.tolerance=all pk.fields=ID auto.evolve=true
异常表现
此前运行完全正常,未做实质性修改的情况下突然报错,测试之前可用的备份配置仍出现问题。排查发现Oracle中ID NOT NULL NUMBER字段被Debezium封装为STRUCT类型,结构如下:
"type":"struct", "fields":[ { "type":"int32", "optional":false, "field":"scale" }, { "type":"bytes", "optional":false, "field":"value" } ],
1. 调整Debezium Oracle连接器的数值类型处理模式
Debezium默认会将未指定精度/标度的Oracle NUMBER类型序列化为STRUCT(确保精度不丢失),但JDBC Sink连接器无法直接处理该类型。在Oracle源连接器配置中添加以下参数,将数值类型直接转换为普通数值:
decimal.handling.mode=numeric
该配置会将NUMBER类型转换为Kafka Connect的NUMERIC类型,JDBC Sink可直接识别。如果允许精度损失,也可设置为double,但numeric能更好保留原始精度。
2. 检查Oracle表结构是否隐性变更
尽管你未修改配置,仍需确认Oracle表的ID字段定义是否被他人修改。执行以下SQL查看表结构:
DESC DEBEZIUM.CUSTOMER;
确认ID字段是否仍为NOT NULL NUMBER,是否新增了精度/标度定义,或者存在其他隐性变更。
3. 排查Debezium版本或依赖变更
检查plugin.path指定的目录中,Debezium Oracle连接器的jar包是否被替换或升级。部分Debezium版本(如1.9+)默认将decimal.handling.mode设为precise(即STRUCT),若之前使用的版本默认值为numeric,版本变更会导致该问题。
4. Sink端转换处理STRUCT类型(备选方案)
若无法修改源连接器配置,可在JDBC Sink连接器中添加转换逻辑,提取STRUCT中的实际数值。例如添加ScriptedTransform转换:
transforms=unwrap,parseId transforms.parseId.type=io.debezium.transforms.ScriptedTransform transforms.parseId.script=groovy: import java.nio.ByteBuffer; if (value.ID != null) { // 根据实际数值类型选择getInt()或getLong() value.ID = ByteBuffer.wrap(value.ID.value).getLong(); } return value;
此方法需根据ID字段的实际数值范围调整getInt()或getLong(),确保数值正确解析。
内容的提问来源于stack exchange,提问作者Muhammad Affan

