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

Flink Kafka Sink复杂数据类型无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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 09:35:15