如何从单条Record生成多条Record?流处理查询方案咨询
问题
假设流Stream A里的记录包含{timestamp, json}两个字段,示例记录如下:
{"row":{"columns":[1678710830000, {"1": "abc", "2": "bcd", "3": "xyz"}]}}
需要把这类记录转换成包含{timestamp, fieldA, fieldB}虚拟字段的多条记录,输出到Stream B(或表),目标记录示例如下:
{"row":{"columns":[1678710830000, "1", "abc"]}} {"row":{"columns":[1678710830000, "2", "bcd"]}} {"row":{"columns":[1678710830000, "3", "xyz"]}}
求对应的查询语句方案。
解决方案
下面针对几种主流流处理引擎给出具体实现:
Apache Flink SQL
核心思路是用UNNEST配合JSON对象键值拆解函数,把原JSON里的每个键值对拆成单独记录:
-- 定义输入流Stream A CREATE TABLE Stream_A ( `timestamp` BIGINT, `json` MAP<STRING, STRING> -- 若Flink版本支持,也可直接用JSON类型 ) WITH ( 'connector' = 'kafka', -- 根据实际数据源调整,比如CDC、文件等 'topic' = 'stream-a-topic', 'properties.bootstrap.servers' = 'localhost:9092', 'format' = 'json' ); -- 定义输出流Stream B CREATE TABLE Stream_B ( `timestamp` BIGINT, fieldA STRING, fieldB STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'stream-b-topic', 'properties.bootstrap.servers' = 'localhost:9092', 'format' = 'json' ); -- 转换并插入数据 INSERT INTO Stream_B SELECT s.`timestamp`, entry.key AS fieldA, entry.value AS fieldB FROM Stream_A s, UNNEST(s.`json`) AS entry(key, value);
如果原json字段是JSON字符串而非MAP类型,可先转成MAP再拆解:
INSERT INTO Stream_B SELECT s.`timestamp`, entry.key AS fieldA, entry.value AS fieldB FROM Stream_A s, UNNEST(JSON_TO_MAP(s.`json`)) AS entry(key, value);
Spark Structured Streaming(SQL/DSL)
DSL方式
// 读取Stream A数据 val streamA = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "localhost:9092") .option("subscribe", "stream-a-topic") .load() .select( // 解析JSON结构,提取timestamp和json字段 from_json(col("value").cast("string"), struct(struct( col("columns")(0).cast("bigint").alias("timestamp"), col("columns")(1).cast("map<string,string>").alias("json") ).alias("row")) ).alias("data") ) .select("data.row.timestamp", "data.row.json") // 拆解JSON键值对并输出到Stream B val streamB = streamA .select( col("timestamp"), // 把MAP转成键值对数组后展开 explode(map_entries(col("json"))).alias("entry") ) .select( col("timestamp"), col("entry.key").alias("fieldA"), col("entry.value").alias("fieldB") ) // 输出到Kafka或其他存储 streamB.writeStream .format("kafka") .option("kafka.bootstrap.servers", "localhost:9092") .option("topic", "stream-b-topic") .option("checkpointLocation", "/tmp/spark-checkpoint") .start() .awaitTermination()
SQL方式
-- 创建临时视图映射Stream A CREATE OR REPLACE TEMP VIEW Stream_A AS SELECT cast(get_json_object(value, '$.row.columns[0]') as bigint) AS `timestamp`, get_json_object(value, '$.row.columns[1]') AS `json_str` FROM kafka WHERE topic = 'stream-a-topic'; -- 转换并插入到Stream B对应的表 INSERT INTO stream_b_table SELECT `timestamp`, entry.key AS fieldA, entry.value AS fieldB FROM Stream_A, LATERAL VIEW explode(map_entries(from_json(`json_str`, 'map<string,string>'))) AS entry;
KSQL(Kafka生态专用SQL)
-- 定义输入流Stream A CREATE STREAM Stream_A ( row STRUCT<columns ARRAY<STRING>> ) WITH ( KAFKA_TOPIC='stream-a-topic', VALUE_FORMAT='JSON' ); -- 转换生成Stream B CREATE STREAM Stream_B AS SELECT CAST(row.columns[1] AS BIGINT) AS `timestamp`, entry->key AS fieldA, entry->value AS fieldB FROM Stream_A, UNNEST( -- 提取JSON对象的所有键,再映射成键值对结构体后展开 TRANSFORM( MAP_KEYS(CAST(row.columns[2] AS MAP<STRING, STRING>)), (k) => STRUCT(key := k, value := CAST(row.columns[2] AS MAP<STRING, STRING>)[k]) ) ) AS t(entry);
内容的提问来源于stack exchange,提问作者Aman Agarwal
相关产品推荐
相关产品推荐

