如何解决基于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来解码原始值:
- 修改Sink连接器的schema配置:
{ "name": "azure-adx-kafka-sink-connector", "config": { // 保留原有配置 "key.converter.schemas.enable": "true", "value.converter.schemas.enable": "true" } }
- 配置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值并转为Inttodatetimefromunixtimemilliseconds:将毫秒时间戳转为ADX datetime类型
方案选择建议
- 若需同步到多个目标系统(如PostgreSQL),优先选方案1:统一输出通用可读格式,所有下游系统无需特殊适配
- 若仅针对ADX同步,且不想修改源数据,可选方案2,但后续新增其他系统时需额外适配
内容的提问来源于stack exchange,提问作者Robsn
相关产品推荐
相关产品推荐

