如何在Apache Flink Table API中原样传输异构类型数组
解决方案
要保留values数组中元素的原始类型(int、decimal、bool等),无需解析转换直接原样复制,可通过以下几种方案实现:
方案1:使用Flink JSON类型(推荐,Flink 1.13+)
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
相关产品推荐
相关产品推荐

