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

如何解决基于Kafka Connect的Oracle到ADX数据类型转换问题

问题

目标是将Oracle数据库的数据导入Azure Data Explorer(ADX),未来可能还要同步至PostgreSQL等其他系统。当前使用Confluent JDBC连接器从Oracle抽取数据,整体运行正常,但ADX写入时非字符串类型数据存在问题:

  • ADX列设为Int类型时,导入值为NULL
  • ADX列设为String类型时,写入的是Kafka中显示的原始编码值(如AMOUNT字段的AtA=)
  • 尝试将Sink连接器的value.converter.schemas.enable设为true后,所有值均为NULL

Connect容器通用配置(YAML)

kafka-connect:
  image: confluentinc/cp-kafka-connect:7.7.2
  container_name: kafka-connect
  depends_on:
    - kafka-0
    - kafka-1
    - kafka-2
  volumes:
    - /home/portainer/kafka-connect:/etc/kafka-connect/jars/
  environment:
    - CONNECT_BOOTSTRAP_SERVERS=kafka-0:9092,kafka-1:9092,kafka-2:9092
    - CONNECT_GROUP_ID=connect-cluster
    - CONNECT_CONFIG_STORAGE_TOPIC=connect-configs
    - CONNECT_OFFSET_STORAGE_TOPIC=connect-offsets
    - CONNECT_STATUS_STORAGE_TOPIC=connect-status
    - CONNECT_KEY_CONVERTER=org.apache.kafka.connect.json.JsonConverter
    - CONNECT_KEY_CONVERTER_SCHEMAS_ENABLE=false  
    - CONNECT_VALUE_CONVERTER=org.apache.kafka.connect.json.JsonConverter
    - CONNECT_PLUGIN_PATH=/usr/share/java/,/etc/kafka-connect/jars/
    - CONNECT_REST_ADVERTISED_HOST_NAME=kafka-connect
    - CONNECT_TOPIC_CREATION_ENABLE=true
    - CONNECT_STATUS_STORAGE_PARTITIONS=5
  ports:
    - "8083:8083"

Oracle源连接器配置(JSON)

{
  "name": "oracle-jdbc-source-connector",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
    "tasks.max": "1",
    "connection.user": "XXX",
    "connection.password": "XXX",
    "connection.url": "jdbc:oracle:thin:@XXX:1521/XXX",
    "table.whitelist": "ORDER-TABLE",
    "mode": "timestamp",
    "timestamp.column.name": "LAST_CHANGE_TS",
    "topic.prefix": "ORACLE-",
    "poll.interval.ms": "10000"
  }
}

Kafka消息示例

{
   "schema": {
      "type": "struct",
      "fields": [
         {
            "type": "string",
            "optional": false,
            "field": "ORDER"
         },
         {
            "type": "bytes",
            "optional": true,
            "name": "org.apache.kafka.connect.data.Decimal",
            "version": 1,
            "parameters": {
               "scale": "0",
               "connect.decimal.precision": "38"
            },
            "field": "AMOUNT"
         },
         {
            "type": "int64",
            "optional": false,
            "name": "org.apache.kafka.connect.data.Timestamp",
            "version": 1,
            "field": "LAST_CHANGE_TS"
         }
      ],
      "optional": false,
      "name": "ORDER-TABLE"
   },
   "payload": {
      "ORDER": "2246573",
      "AMOUNT": "AtA=",
      "LAST_CHANGE_TS": 1677153916000
   }
}

ADX Sink连接器配置(JSON)

{
    "name": "azure-adx-kafka-sink-connector",
    "config": {
        "connector.class": "com.microsoft.azure.kusto.kafka.connect.sink.KustoSinkConnector",
        "flush.size.bytes": 1000,
        "flush.interval.ms": 1000,
        "tasks.max": 1,
        "topics": "ORACLE-ORDER-TABLE",
        "kusto.tables.topics.mapping": "[{'topic': 'ORACLE-ORDER-TABLE', 'db': 'TEST', 'table': 'ORDER_T', 'format': 'json', 'mapping':'ORDER_T_Mapping'}]",
        "aad.auth.authority": "XXX",
        "aad.auth.appid": "XXX",
        "aad.auth.appkey": "XXX",
        "kusto.ingestion.url": "XXX",
        "kusto.query.url": "XXX",
        "key.converter": "org.apache.kafka.connect.json.JsonConverter",
        "value.converter": "org.apache.kafka.connect.json.JsonConverter",
        "key.converter.schemas.enable": "false",
        "value.converter.schemas.enable": "false"
    }
}

解决方案

核心问题分析

Kafka消息中的Decimal、Timestamp类型被编码为Base64字节、毫秒时间戳等非直接可读格式,ADX Sink连接器关闭schema解析时无法识别这些类型;开启schema解析后,因配置与消息schema结构不匹配,导致值提取失败。

方案1:修改源连接器输出通用可读格式(推荐,兼容多目标系统)

在Oracle源连接器中添加配置,将特殊类型转换为所有下游系统都能识别的原始格式,无需修改Sink配置,同时兼容未来同步到PostgreSQL的需求:

{
  "name": "oracle-jdbc-source-connector",
  "config": {
    // 保留原有配置
    "decimal.handling.mode": "string",
    "timestamp.converter": "org.apache.kafka.connect.storage.StringConverter",
    "timestamp.column.name": "LAST_CHANGE_TS"
  }
}
  • decimal.handling.mode: string:将Oracle Decimal类型转为原始数值的字符串形式,避免Base64编码,ADX和PostgreSQL都能直接转换为数值类型
  • timestamp.converter:将毫秒时间戳转为ISO标准字符串(如2023-02-23T12:45:16.000+0800),所有主流数据库都能自动解析为时间类型

修改后Kafka消息的payload会变成:

{
  "ORDER": "2246573",
  "AMOUNT": "1234",
  "LAST_CHANGE_TS": "2023-02-23T12:45:16.000+0800"
}

方案2:调整ADX Sink与映射规则适配原始编码格式

如果不想修改源数据格式,可开启Sink的schema解析,并调整ADX ingestion mapping来解码原始值:

  1. 修改Sink连接器的schema配置:
{
    "name": "azure-adx-kafka-sink-connector",
    "config": {
        // 保留原有配置
        "key.converter.schemas.enable": "true",
        "value.converter.schemas.enable": "true"
    }
}
  1. 配置ADX的ingestion mapping,添加转换逻辑:
// 创建表(若未创建)
.create table ORDER_T (ORDER string, AMOUNT int, LAST_CHANGE_TS datetime)

// 创建映射规则
.create table ORDER_T ingestion json mapping 'ORDER_T_Mapping' '['
'{"column":"ORDER", "path":"$.payload.ORDER", "datatype":"string"},'
'{"column":"AMOUNT", "path":"$.payload.AMOUNT", "datatype":"int", "transform":"base64_decode_string"},'
'{"column":"LAST_CHANGE_TS", "path":"$.payload.LAST_CHANGE_TS", "datatype":"datetime", "transform":"todatetimefromunixtimemilliseconds"}'
']'
  • base64_decode_string:解码AMOUNT字段的Base64值并转为Int
  • todatetimefromunixtimemilliseconds:将毫秒时间戳转为ADX datetime类型

方案选择建议

  • 若需同步到多个目标系统(如PostgreSQL),优先选方案1:统一输出通用可读格式,所有下游系统无需特殊适配
  • 若仅针对ADX同步,且不想修改源数据,可选方案2,但后续新增其他系统时需额外适配

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 05:05:53