如何配置Kafka JDBC Connector将Oracle CLOB中JSON转为JSON对象
问题
我在Oracle 11g数据库的INTERFACE_KAFKA_QUEUE表中,将格式合法的JSON数据存储在MESSAGE字段(CLOB类型)中。希望通过Kafka JDBC Connector将该JSON数据拉取并发送至Kafka主题,但目前主题中的MESSAGE值是字符串形式而非JSON对象。
当前主题输出格式
"MESSAGE": "{\"BatchNumber\":\"861621\",\"InterchangeCount\":40,\"ActionType\":\"C\"}"
期望格式
"MESSAGE":
{"BatchNumber":"861621","InterchangeCount":40,"ActionType":"C"}
当前连接器配置
name = [MY-CONNECTOR-NAME] connector.class = io.confluent.connect.jdbc.JdbcSourceConnector key.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=false errors.log.enable = true errors.log.include.messages = true connection.url = [MYDB-CONNECTION-STRING] connection.user = [MYDB-USER] connection.password = [MYDB-PASSWORD] connection.attempts = 2 connection.backoff.ms = 10000 numeric.mapping = best_fit dialect.name = OracleDatabaseDialect mode = incrementing incrementing.column.name = AUDIT_NUMBER query = [MY-SQL-STRING] topic.prefix = [MY-TOPIC] principal.service.name = [MY-SERVICE-NAME] principal.service.password = [MY-SERVICE-PASSWORD]
请问是否可通过配置连接器实现需求?若可以,需如何设置?我不想新增额外组件或使用KSQL。
解决方案
可以通过调整SQL查询和连接器配置实现需求,无需额外组件或KSQL,具体操作如下:
修改SQL查询语句
Oracle 11g的CLOB字段存储的JSON会被连接器默认识别为字符串,需要在查询时将其解析为JSON结构。使用Oracle自带的DBMS_JSON包函数处理:SELECT AUDIT_NUMBER, DBMS_JSON.GET_JSON_OBJECT(MESSAGE) AS MESSAGE FROM INTERFACE_KAFKA_QUEUE若存储的JSON存在转义字符,可先清理转义再解析:
SELECT AUDIT_NUMBER, DBMS_JSON.GET_JSON_OBJECT(REPLACE(MESSAGE, '\\"', '"')) AS MESSAGE FROM INTERFACE_KAFKA_QUEUE添加连接器转换配置
在现有配置中加入以下转换参数,让连接器将解析后的JSON字段直接序列化为JSON对象:transforms=unwrapJson transforms.unwrapJson.type=org.apache.kafka.connect.transforms.UnwrapFromJson$Value transforms.unwrapJson.schemas.enable=false确认转换器配置
你当前已设置value.converter=org.apache.kafka.connect.json.JsonConverter且value.converter.schemas.enable=false,此配置可确保连接器将识别到的JSON结构直接输出为Kafka主题中的JSON对象,而非字符串形式。
内容的提问来源于stack exchange,提问作者Tobie van der Merwe
相关产品推荐
相关产品推荐

