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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 04:05:18