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

通过AWS Glue将Snowflake数据写入DynamoDB遇验证错误求助

问题解决:AWS Glue从Snowflake写入DynamoDB报ValidationException

问题背景

从Snowflake查询得到的DynamicFrame结构如下(字段名带双引号):

root
|-- "customer_id": string
|-- "non_personalized": string
|-- "emalta1": string
|-- "emalta2": string
|-- "emalta3": string

尝试三种写入DynamoDB的方式均报错:

An error occurred while calling o122.pyWriteDynamicFrame. The provided key element does not match the schema (Service: AmazonDynamoDBv2; Status Code: 400; Error Code: ValidationException; Request ID: L7ASM29EH86UQBMKLL5KUF61VBVV4KQNSO5AEMVJF66Q9ASUAAJG; Proxy: null)

核心原因

Snowflake查询中使用了带双引号的标识符(如"db_dev_data"),导致Glue读取的DynamicFrame字段名被保留了双引号,而DynamoDB表的主键是无引号的customer_id,字段名不匹配导致主键验证失败。

解决方案步骤

1. 清理字段名(去掉双引号)

使用Glue的RenameField转换,将带引号的字段名改为无引号的名称,确保与DynamoDB表的字段名完全一致。

2. 转换数据结构(可选,按需处理非主键字段)

将emalta1/emalta2/emalta3合并为数组或结构体,适配DynamoDB的非关系型存储模式。

3. 写入DynamoDB

使用修正后的DynamicFrame执行写入操作。

修改后的完整代码

import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
from awsglue.dynamicframe import DynamicFrame

## @params: [JOB_NAME]
args = getResolvedOptions(sys.argv, ['JOB_NAME'])

sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)

# 从Snowflake读取数据
connection_type="snowflake"
connection_options={
        "connectionName": "teste-snowflake",
        "sfDatabase": '"db_dev_data"',
        "sfSchema": '"tmp"',
        "sfWarehouse": "DBT_WH",
        "query": 'select * from "db_dev_data"."tmp"."fake_data" limit 200'
}
snowflake_read = glueContext.create_dynamic_frame.from_options(connection_type ,connection_options)

print('原始Snowflake DF schema')
print(snowflake_read.printSchema())

# 步骤1:重命名字段,去掉双引号
renamed_df = RenameField.apply(
    frame=snowflake_read,
    old_name='"customer_id"',
    new_name='customer_id'
)
renamed_df = RenameField.apply(frame=renamed_df, old_name='"non_personalized"', new_name='non_personalized')
renamed_df = RenameField.apply(frame=renamed_df, old_name='"emalta1"', new_name='emalta1')
renamed_df = RenameField.apply(frame=renamed_df, old_name='"emalta2"', new_name='emalta2')
renamed_df = RenameField.apply(frame=renamed_df, old_name='"emalta3"', new_name='emalta3')

print('重命名后的DF schema')
print(renamed_df.printSchema())

# 步骤2:将emalta字段合并为数组(可选,按需调整)
transformed_df = ApplyMapping.apply(
    frame=renamed_df,
    mappings=[
        ("customer_id", "string", "customer_id", "string"),
        ("non_personalized", "string", "non_personalized", "string"),
        ("emalta1", "string", "emalta1", "string"),
        ("emalta2", "string", "emalta2", "string"),
        ("emalta3", "string", "emalta3", "string")
    ]
)

# 使用Spark SQL创建数组字段(也可以用Glue自定义转换)
spark_df = transformed_df.toDF()
spark_df = spark_df.withColumn("emaltas", spark_df.array(spark_df.emalta1, spark_df.emalta2, spark_df.emalta3))
spark_df = spark_df.drop("emalta1", "emalta2", "emalta3")

# 转回DynamicFrame
final_df = DynamicFrame.fromDF(spark_df, glueContext, "final_df")

print('最终DF schema')
print(final_df.printSchema())

# 步骤3:写入DynamoDB
glueContext.write_dynamic_frame_from_options(
    frame=final_df,
    connection_type="dynamodb",
    connection_options={
        "dynamodb.output.tableName": "teste-glue-snowflake",
        # 可选:控制写入吞吐量占比
        "dynamodb.throughput.write.percent": "1.0"
    }
)

print("写入完成")
job.commit()

额外注意事项

  • 确保DynamoDB表的主键(customer_id)类型为字符串,与DynamicFrame中的字段类型完全匹配。
  • 若不需要合并字段,可直接跳过步骤2,使用重命名后的DynamicFrame写入。
  • 若Snowflake查询时能避免使用带引号的标识符(如改用小写或符合Spark命名规则的名称),可从源头避免字段名带引号的问题。

内容的提问来源于stack exchange,提问作者Geo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 16:34:53