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

按Key合并DataFrame中JSON字符串为指定JSON数组的技术问询

问题场景

我有一个包含key列和JSON字符串列jsonString的DataFrame,数据如下:

keyjsonString
111{"id" : "12345", "foo" : "stuff"}
111{"id" : "23456", "bar" : "other stuff"}
111{"id" : "34567", "baz" : "even other stuff"}

需求是:针对每个key,将对应的所有JSON字符串合并为一个JSON数组,并添加固定字段type: "combined",最终输出格式如下(用于发布到Kafka Topic):

{
    "type" : "combined",
    "values" :
    [
        {"id" : "12345", "foo" : "stuff"},
        {"id" : "23456", "bar" : "other stuff"},
        {"id" : "34567", "baz" : "even other stuff"}
    ]
}

尝试过直接拼接字符串效果很差,想找无需为所有JSON结构构建巨型Schema的实现方式(实际有4种不同的JSON Schema)。

解决方案

不用提前定义复杂Schema,直接在字符串层面操作即可完成需求,以PySpark为例:

  1. 按key分组,收集JSON字符串数组
    先对DataFrame按key分组,把每组的jsonString收集成字符串数组:

    from pyspark.sql import functions as F
    
    grouped_df = df.groupBy("key").agg(F.collect_list("jsonString").alias("raw_values"))
    
  2. 将字符串数组转为标准JSON数组格式
    用concat_ws把数组里的JSON字符串用逗号连接,再包裹数组前后的方括号:

    formatted_df = grouped_df.withColumn(
        "values",
        F.concat(F.lit("["), F.concat_ws(", ", "raw_values"), F.lit("]"))
    )
    
  3. 拼接生成最终Kafka消息JSON
    把固定字段type和生成的values数组拼接成完整的JSON结构:

    final_df = formatted_df.withColumn(
        "kafka_message",
        F.concat(
            F.lit('{"type": "combined", "values": '),
            "values",
            F.lit("}")
        )
    ).drop("raw_values", "values", "key")
    
  4. 发布到Kafka Topic
    直接将kafka_message列作为消息内容发布:

    final_df.selectExpr("CAST(kafka_message AS STRING) AS value") \
            .write \
            .format("kafka") \
            .option("kafka.bootstrap.servers", "your_broker:port") \
            .option("topic", "your_target_topic") \
            .save()
    

补充说明

  • 这种方式完全不需要解析JSON结构,适配任意JSON Schema的场景;
  • 如果JSON字符串存在换行或特殊字符,可先用F.regexp_replace清洗,比如F.regexp_replace("jsonString", "\n", "");
  • 若使用Scala,逻辑完全一致,仅函数调用语法略有差异。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 19:33:31