通过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
相关产品推荐
相关产品推荐

