如何基于DataFrame(DynamicFrame)其他列的JSON数据添加新列
最优实现方案
优先选择Spark原生算子实现,性能远高于自定义UDF或行级遍历方案,以下分两种常用场景给出代码:
场景1:使用PySpark DataFrame
如果customJson字段是JSON字符串类型,先定义JSON结构schema再解析提取即可:
from pyspark.sql import functions as F from pyspark.sql.types import StructType, StringType # 定义customJson对应的结构体schema json_schema = StructType() \ .add("key", StringType()) \ .add("value", StringType()) # 解析JSON并提取value作为lastName列 result_df = df.withColumn("json_parsed", F.from_json(F.col("customJson"), json_schema)) \ .withColumn("lastName", F.col("json_parsed.value")) \ .drop("json_parsed") # 可选删除中间临时列
如果customJson本身已经是Spark结构体类型而非字符串,可跳过解析步骤直接提取:
result_df = df.withColumn("lastName", F.col("customJson.value"))
如果需要根据JSON中key的值动态生成列名(不同行key值不同的场景),可以配合pivot实现:
result_df = df.withColumn("json_parsed", F.from_json(F.col("customJson"), json_schema)) \ .select("id", "name", "customJson", "json_parsed.*") \ .groupBy("id", "name", "customJson") \ .pivot("key") \ .agg(F.first("value"))
场景2:使用AWS Glue DynamicFrame
如果使用的是Glue DynamicFrame,推荐先转成DataFrame用上述方法处理后转回,性能最优:
from awsglue.dynamicframe import DynamicFrame # DynamicFrame转DataFrame处理 df = dynamic_frame.toDF() # 调用上述DataFrame处理逻辑得到result_df processed_dyf = DynamicFrame.fromDF(result_df, glueContext, "processed_dyf")
也可以用Glue原生的Map转换实现行级处理,适合逻辑非常复杂的场景:
def extract_json_field(rec): import json json_data = json.loads(rec["customJson"]) rec[json_data["key"]] = json_data["value"] return rec processed_dyf = Map.apply(frame = dynamic_frame, f = extract_json_field)
内容的提问来源于stack exchange,提问作者Miroslav Petrovic
相关产品推荐
相关产品推荐

