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

如何在Apache Flink Table API中原样传输异构类型数组

解决方案

要保留values数组中元素的原始类型(int、decimal、bool等),无需解析转换直接原样复制,可通过以下几种方案实现:

Flink 1.13及以上版本支持原生JSON类型,可直接存储任意JSON结构并保留原始类型信息。只需修改表定义中values字段的类型为JSON:

源表定义

CREATE TEMPORARY TABLE data_topic (
  `hash` STRING,
  `date` STRING,
  `values` JSON, -- 用JSON类型存储混合类型数组
  data_event_time AS TO_TIMESTAMP(`date`, 'yyyyMMdd''T''HHmmss.SSSX'),
  WATERMARK FOR data_event_time AS data_event_time
) WITH (
  'connector' = 'kafka',
  'topic' = 'name_of_topic',
  'scan.startup.mode' = 'latest-offset',
  'key.format' = 'json',
  'key.fields' = 'hash',
  'value.format' = 'json',
  'value.fields-include' = 'EXCEPT_KEY'
);

Sink表定义

Sink表的values字段同样定义为JSON类型,确保输出时原样序列化:

CREATE TEMPORARY TABLE sink_topic (
  `hash` STRING,
  `date` STRING,
  `values` JSON
) WITH (
  'connector' = 'kafka',
  'topic' = 'name_of_sink_topic',
  'key.format' = 'json',
  'key.fields' = 'hash',
  'value.format' = 'json',
  'value.fields-include' = 'EXCEPT_KEY'
);

同步语句

直接字段映射即可:

INSERT INTO sink_topic SELECT hash, date, values FROM data_topic;

这种方式无需额外依赖,Flink会自动保留values数组的原始类型结构,是最简单的实现方式。

方案2:使用RAW类型(适配低版本Flink)

如果使用Flink 1.13以下版本,可通过RAW类型存储原始JSON结构,配合Jackson序列化器实现类型保留:

源表定义

CREATE TEMPORARY TABLE data_topic (
  `hash` STRING,
  `date` STRING,
  -- 用RAW类型存储JsonNode,指定Jackson序列化器
  `values` RAW('com.fasterxml.jackson.databind.JsonNode', 'org.apache.flink.table.data.util.JsonRawValueSerializer'),
  data_event_time AS TO_TIMESTAMP(`date`, 'yyyyMMdd''T''HHmmss.SSSX'),
  WATERMARK FOR data_event_time AS data_event_time
) WITH (
  'connector' = 'kafka',
  'topic' = 'name_of_topic',
  'scan.startup.mode' = 'latest-offset',
  'key.format' = 'json',
  'key.fields' = 'hash',
  'value.format' = 'json',
  'value.fields-include' = 'EXCEPT_KEY'
);

Sink表定义

Sink表的values字段需保持一致的RAW类型定义,确保输出时原样序列化。

方案3:直接处理完整JSON Value

如果需要更灵活的原始数据控制,可以将整个Kafka Value当作JSON对象,手动提取所需字段:

源表定义

CREATE TEMPORARY TABLE data_topic (
  `key` STRING,
  `value` JSON, -- 完整Kafka Value作为JSON存储
  -- 通过JSON函数提取字段,values保留原始数组结构
  `hash` STRING AS JSON_VALUE(`value`, '$.hash'),
  `date` STRING AS JSON_VALUE(`value`, '$.date'),
  `values` JSON AS JSON_QUERY(`value`, '$.values'),
  data_event_time AS TO_TIMESTAMP(`date`, 'yyyyMMdd''T''HHmmss.SSSX'),
  WATERMARK FOR data_event_time AS data_event_time
) WITH (
  'connector' = 'kafka',
  'topic' = 'name_of_topic',
  'scan.startup.mode' = 'latest-offset',
  'key.format' = 'json',
  'key.fields' = 'hash',
  'value.format' = 'raw',
  'value.raw.format' = 'utf-8'
);

这种方式适合需要频繁调整提取字段的场景,无需修改表结构即可扩展字段。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 00:23:13