Flink Kafka Sink复杂数据类型无UDF实现方案咨询
Flink Kafka Connector 处理复杂嵌套JSON输出的无UDF方案
问题场景
从Kafka的SOURCE_TABLE读取数据,其中contentJson是嵌套JSON字符串,示例:
"contentJson": "{\"field1\":[{\"identityId\":\"21ade8be05f548a983acf68a8ab5ff73\",\"IdsAfter\":[\"seg-id-1\"]}]}"
需要将field1输出为JSON数组(无引号包裹),但不想编写UDF,也不想在Sink表中定义ARRAY<ROW<...>>这类复杂类型,尝试RAW类型报错。
解决方案:利用JSON格式的raw-mode参数
通过Flink内置JSON函数提取目标片段,结合Sink表的raw-mode配置,直接输出正确格式的JSON。
1. 定义SOURCE_TABLE
正常读取Kafka的JSON数据,contentJson设为STRING类型:
CREATE TABLE SOURCE_TABLE ( contentJson STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'your-source-topic', 'properties.bootstrap.servers' = 'xxx:9092', 'format' = 'json' );
2. 定义SINK_TABLE
field1设为STRING类型,开启json.raw-mode参数(核心配置):
CREATE TABLE SINK_TABLE ( field1 STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'your-sink-topic', 'properties.bootstrap.servers' = 'xxx:9092', 'format' = 'json', 'json.raw-mode' = 'true' -- 开启后STRING字段直接输出原始内容,不添加引号 );
3. 编写插入逻辑
用内置JSON_QUERY函数提取contentJson中的field1片段,无需UDF:
INSERT INTO SINK_TABLE SELECT -- 提取field1的原始JSON片段,异常时返回null JSON_QUERY(contentJson, '$.field1' NULL ON ERROR) AS field1 FROM SOURCE_TABLE;
关键说明
JSON_QUERY:从嵌套JSON字符串中提取指定路径的JSON片段,返回未转义的JSON内容(如[{"identityId":"...","IdsAfter":["..."]}])。json.raw-mode:开启后,Flink的JSON格式会将STRING类型字段的原始内容直接输出为JSON值,不会包裹引号,正好匹配需求。- RAW类型报错原因:RAW类型需要指定特定序列化器,且输出为原始字节,与JSON格式的序列化逻辑不兼容,因此不适用该场景。
内容的提问来源于stack exchange,提问作者hitesh
相关产品推荐
相关产品推荐

