Spark 2.1.0:PySpark DataFrame列转嵌套JSON格式的实现问询
PySpark 2.1.0: Transform DataFrame to tag + JSON-structured data column
我来帮你搞定这个DataFrame转换需求,在Spark 2.1.0版本里,我们可以通过组合create_map和to_json函数来构造嵌套的JSON结构,完全匹配你的预期输出。下面是完整的实现代码:
from pyspark.sql import SparkSession from pyspark.sql.functions import create_map, lit, col, to_json spark = SparkSession.builder.appName("example").getOrCreate() df = spark.createDataFrame([ (1388534400, "GOOG", 50, 'a', 1), (1388534400, "FB", 60, 'b', 2), (1388534400, "MSFT", 55, 'c', 3), (1388620800, "GOOG", 52, 'd', 4)] ).toDF("date", "stock", "price", 'tag', 'num') # 构造嵌套Map并转换为JSON字符串 result_df = df.select( col("tag"), to_json( create_map( lit("A"), create_map(lit("stock"), col("stock"), lit("price"), col("price")), lit("B"), create_map(lit("date"), col("date"), lit("num"), col("num")) ) ).alias("data") ) # 查看最终结果 result_df.show(truncate=False)
关键步骤解释:
- 构造内层映射:用
create_map分别为A、B字段生成键值对结构,比如create_map(lit("stock"), col("stock"), lit("price"), col("price"))会生成{'stock': 'GOOG', 'price': 50}这样的Map - 组合外层嵌套:再用一层
create_map把A和B作为顶级键,拼接成完整的嵌套Map结构 - 转JSON字符串:通过
to_json函数把嵌套Map转换成标准JSON格式的字符串,命名为data列 - 保留目标列:最终只选取
tag和data两列作为结果
额外说明:
Spark生成的JSON会使用标准双引号(符合JSON规范),如果你的场景确实需要单引号格式,可以额外添加regexp_replace做替换:
from pyspark.sql.functions import regexp_replace result_df = result_df.withColumn( "data", regexp_replace(col("data"), '"', "'") )
内容的提问来源于stack exchange,提问作者wffzxyl
相关产品推荐
相关产品推荐

