You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.03 19:05:20