按Key合并DataFrame中JSON字符串为指定JSON数组的技术问询
问题场景
我有一个包含key列和JSON字符串列jsonString的DataFrame,数据如下:
| key | jsonString |
|---|---|
| 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为例:
按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"))将字符串数组转为标准JSON数组格式
用concat_ws把数组里的JSON字符串用逗号连接,再包裹数组前后的方括号:formatted_df = grouped_df.withColumn( "values", F.concat(F.lit("["), F.concat_ws(", ", "raw_values"), F.lit("]")) )拼接生成最终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")发布到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
相关产品推荐
相关产品推荐

