如何通过Glue写入DynamoDB时避免AttributeValues,输出原生数组
问题描述
我有一个Schema如下的DynamicFrame:
root |-- data1: string (nullable = false) |-- data2: string (nullable = false) |-- data3: array (nullable = false) | |-- element: string (containsNull = true)
使用以下代码将其写入DynamoDB时:
glue_context.write_dynamic_frame_from_options( frame=DynamicFrame.fromDF(df, glue_context, "output"), connection_type="dynamodb", connection_options={ "dynamodb.output.tableName": "table_name", "dynamodb.throughput.write.percent": "1.0", }, )
data3字段被写入为[ { "L" : [ { "S" : "" }, { "S" : "" }, { "S" : "" }, { "S" : "" } ] } ],但我希望它以["","","",""]的格式输出,该如何实现?
解决方案
方法1:添加DynamoDB序列化选项
Glue的DynamoDB连接器默认会给复杂类型加上类型标记结构,你可以通过添加序列化选项强制使用原生类型映射:
dynamic_frame = DynamicFrame.fromDF(df, glue_context, "output") glue_context.write_dynamic_frame_from_options( frame=dynamic_frame, connection_type="dynamodb", connection_options={ "dynamodb.output.tableName": "table_name", "dynamodb.throughput.write.percent": "1.0", # 启用原生DynamoDB类型序列化 "dynamodb.serde": "org.apache.hadoop.dynamodb.DynamoDBItemWritable" }, )
这个配置会让Spark字符串数组直接映射为DynamoDB的L类型列表,后续用SDK读取时会自动解析为普通字符串数组["","","",""]。
方法2:显式确认字段类型
如果字段类型被意外转换,先强制指定data3为字符串数组类型再写入:
from pyspark.sql.functions import col # 显式转换data3为字符串数组 processed_df = df.withColumn("data3", col("data3").cast("array<string>")) dynamic_frame = DynamicFrame.fromDF(processed_df, glue_context, "processed_output") glue_context.write_dynamic_frame_from_options( frame=dynamic_frame, connection_type="dynamodb", connection_options={ "dynamodb.output.tableName": "table_name", "dynamodb.throughput.write.percent": "1.0", }, )
方法3:存储为JSON字符串(若需要纯JSON格式)
如果希望在DynamoDB中直接存储["","","",""]格式的JSON字符串(而非DynamoDB原生列表类型),可以将数组序列化为JSON字符串:
from pyspark.sql.functions import to_json, col # 将data3数组转为JSON字符串 processed_df = df.withColumn("data3", to_json(col("data3"))) dynamic_frame = DynamicFrame.fromDF(processed_df, glue_context, "json_output") glue_context.write_dynamic_frame_from_options( frame=dynamic_frame, connection_type="dynamodb", connection_options={ "dynamodb.output.tableName": "table_name", "dynamodb.throughput.write.percent": "1.0", }, )
内容的提问来源于stack exchange,提问作者Sundeep
相关产品推荐
相关产品推荐

