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

如何从单条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"]}}

求对应的查询语句方案。

解决方案

下面针对几种主流流处理引擎给出具体实现:

核心思路是用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 12:53:26