pySpark如何将DataFrame嵌套Map列转标准JSON字符串插入DynamoDB
问题原因
直接将Struct/Map类型字段通过cast(StringType())转换为字符串时,Spark只会拼接字段值,不会保留键名和JSON结构,因此会出现你遇到的无键名、格式混乱的问题。
解决方法
修改批量转换列类型的逻辑,对复杂类型字段使用to_json函数做序列化,普通字段保留原有的强制类型转换逻辑即可,修改后的代码如下:
import pyspark.sql.types as T from pyspark.sql import functions as SF df = spark.read.option("multiline", "true").json('/home/abhishek.tirkey/Documents/test') Records = df.withColumn("Records", SF.explode(SF.col("Records"))) Rows = Records.select( "Records.column1", "Records.column2", "Records.column3", "Records.column4", ) # 修改列转换逻辑 for col_name in Rows.columns: col_type = Rows.schema[col_name].dataType # 结构化类型、Map类型用to_json转成标准JSON字符串 if isinstance(col_type, (T.StructType, T.MapType)): Rows = Rows.withColumn(col_name, SF.to_json(SF.col(col_name))) # 基础类型直接转字符串 else: Rows = Rows.withColumn(col_name, Rows[col_name].cast(T.StringType())) RowsJSON = Rows.toJSON()
效果验证
转换完成后可执行Rows.select("column4").show(truncate=False)查看输出,此时column4已经是符合要求的标准JSON字符串,直接插入DynamoDB即可得到预期格式。
内容的提问来源于stack exchange,提问作者Abhishek Tirkey
相关产品推荐
相关产品推荐

