Kafka JDBC源连接器(Oracle):如何正确映射嵌套JSON到Avro Schema
问题描述
我使用Confluent的Kafka JDBC Source Connector从Oracle数据库抽取数据推送至Kafka Topic。数据库表test_table包含json_obj(VARCHAR2类型,存储嵌套JSON)和last_update_date(DATE类型)两列。
json_obj的嵌套JSON结构示例:
{ "name": "John", "department": { "deptId": "1", "deptName": "dept1" } }
目标Kafka Topic的Avro Schema定义:
{ "fields": [ { "name": "name", "type": "string" }, { "name": "department", "type":[ { "fields": [ { "name": "deptId", "type": "string" }, { "name": "deptName", "type": [ "string", "null" ] } ], "name": "DeptObj", "type": "record" },"null"] } ], "name": "EmpObj", "namespace": "com.test", "type": "record" }
我的连接器配置:
{ "name": "JdbcSourceConnectorConnector_0", "config": { "schema.registry.url": "<schema-reg-url>", "name": "JdbcSourceConnectorConnector_0", "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector", "tasks.max": "1", "connection.url": "jdbc:oracle:thin:@<ldap-connect-string>/<db-name>", "connection.user": "<db-schema-name>", "connection.password": "<db-pwd>", "numeric.mapping": "best_fit", "dialect.name": "OracleDatabaseDialect", "mode": "timestamp", "timestamp.column.name": "last_update_date", "validate.non.null": "false", "query": "SELECT * FROM (SELECT JSON_VALUE(json_obj, '$.empId') as \"empId\", JSON_VALUE(json_obj, '$.department.deptId') as \"department.deptId\", JSON_VALUE(json_obj, '$.department.deptName') as \"department.deptName\" FROM test_table) A", "table.types": "TABLE", "poll.interval.ms": "30000", "topic.prefix": "<my-topic>", "db.timezone": "America/Los_Angeles", "key.serializer": "org.apache.kafka.common.serialization.StringSerializer", "value.serializer": "io.confluent.kafka.serializers.KafkaAvroSerializer", "transforms": "createKeyStruct,ExtractField, addNamespace", "transforms.createKeyStruct.fields": "empId", "transforms.createKeyStruct.type": "org.apache.kafka.connect.transforms.ValueToKey", "transforms.ExtractField.field": "empId", "transforms.ExtractField.type": "org.apache.kafka.connect.transforms.ExtractField$Key", "transforms.addNamespace.type":"org.apache.kafka.connect.transforms.SetSchemaMetadata$Value", "transforms.addNamespace.schema.name": "EmpObj" } }
运行时抛出错误:
Caused by: org.apache.avro.SchemaParseException: Illegal character in: department.deptId at org.apache.avro.Schema.validateName(Schema.java:1561) at org.apache.avro.Schema.access$400(Schema.java:87) at org.apache.avro.Schema$Field.<init>(Schema.java:541) at org.apache.avro.Schema$Field.<init>(Schema.java:580) at io.confluent.connect.avro.AvroData.addAvroRecordField(AvroData.java:1114) at io.confluent.connect.avro.AvroData.fromConnectSchema(AvroData.java:910) at io.confluent.connect.avro.AvroData.fromConnectSchema(AvroData.java:732) at io.confluent.connect.avro.AvroData.fromConnectSchema(AvroData.java:726) at io.confluent.connect.avro.AvroConverter.fromConnectData(AvroConverter.java:85) at org.apache.kafka.connect.storage.Converter.fromConnectData(Converter.java:63) at org.apache.kafka.connect.runtime.WorkerSourceTask.lambda$convertTransformedRecord$3(WorkerSourceTask.java:321) 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 [2022-10-10 07:53:57,012] INFO Stopping JDBC source task (io.confluent.connect.jdbc.source.JdbcSourceTask)
需要解决的问题:如何正确将department.deptId这类嵌套字段映射到目标Avro Schema?
解决方案
错误核心原因是Avro不允许字段名包含.,你当前SQL中将字段别名设为department.deptId,JDBC连接器会把它当作单个扁平字段名,导致Avro解析失败。要实现嵌套结构映射,可通过以下两种方案解决:
方案1:通过SQL提取扁平字段 + Connect Transforms重组嵌套结构
步骤1:修改SQL查询,用简单别名提取字段
将SQL中的别名改为无.的名称,生成扁平字段:
SELECT JSON_VALUE(json_obj, '$.name') as "name", JSON_VALUE(json_obj, '$.department.deptId') as "deptId", JSON_VALUE(json_obj, '$.department.deptName') as "deptName", last_update_date FROM test_table
步骤2:更新连接器配置,添加Transforms重组嵌套结构
在原有配置基础上,新增Transforms将扁平的deptId、deptName合并为嵌套的department结构,同时移除冗余字段:
{ "name": "JdbcSourceConnectorConnector_0", "config": { "schema.registry.url": "<schema-reg-url>", "name": "JdbcSourceConnectorConnector_0", "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector", "tasks.max": "1", "connection.url": "jdbc:oracle:thin:@<ldap-connect-string>/<db-name>", "connection.user": "<db-schema-name>", "connection.password": "<db-pwd>", "numeric.mapping": "best_fit", "dialect.name": "OracleDatabaseDialect", "mode": "timestamp", "timestamp.column.name": "last_update_date", "validate.non.null": "false", "query": "SELECT JSON_VALUE(json_obj, '$.name') as \"name\", JSON_VALUE(json_obj, '$.department.deptId') as \"deptId\", JSON_VALUE(json_obj, '$.department.deptName') as \"deptName\", last_update_date FROM test_table", "table.types": "TABLE", "poll.interval.ms": "30000", "topic.prefix": "<my-topic>", "db.timezone": "America/Los_Angeles", "key.serializer": "org.apache.kafka.common.serialization.StringSerializer", "value.serializer": "io.confluent.kafka.serializers.KafkaAvroSerializer", "transforms": "createKeyStruct,ExtractField, wrapDept, removeFlatDeptFields, addNamespace", // 原有Key相关转换逻辑 "transforms.createKeyStruct.fields": "empId", "transforms.createKeyStruct.type": "org.apache.kafka.connect.transforms.ValueToKey", "transforms.ExtractField.field": "empId", "transforms.ExtractField.type": "org.apache.kafka.connect.transforms.ExtractField$Key", // 新增嵌套结构转换 "transforms.wrapDept.type": "org.apache.kafka.connect.transforms.HoistField$Value", "transforms.wrapDept.field": "department", "transforms.wrapDept.fields": "deptId,deptName", // 移除扁平的dept字段 "transforms.removeFlatDeptFields.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.removeFlatDeptFields.blacklist": "deptId,deptName", // 指定完整的Schema命名空间 "transforms.addNamespace.type":"org.apache.kafka.connect.transforms.SetSchemaMetadata$Value", "transforms.addNamespace.schema.name": "com.test.EmpObj" } }
Transforms作用说明
wrapDept:将deptId和deptName字段包装到名为department的结构体中removeFlatDeptFields:移除原有的扁平deptId、deptName字段,仅保留嵌套结构addNamespace:必须指定完整命名空间com.test.EmpObj,确保Schema Registry能匹配目标Avro Schema
方案2:使用Oracle JSON函数直接生成嵌套JSON(Oracle 12c+支持)
利用Oracle的JSON_OBJECT函数直接生成嵌套结构的JSON列,让JDBC连接器解析为嵌套的Connect结构:
SELECT JSON_VALUE(json_obj, '$.name') as "name", JSON_OBJECT( 'deptId' VALUE JSON_VALUE(json_obj, '$.department.deptId'), 'deptName' VALUE JSON_VALUE(json_obj, '$.department.deptName') ) as "department", last_update_date FROM test_table
此方案无需额外Transforms,连接器会自动将JSON列解析为嵌套结构,直接匹配目标Avro Schema。
内容的提问来源于stack exchange,提问作者AwsBeginner
相关产品推荐
相关产品推荐

