AWS Glue Job报错:调用o110.pyWriteDynamicFrame时awaitResult抛出异常
AWS Glue写入Redshift报错:An error occurred while calling o110.pyWriteDynamicFrame. Exception is thrown in awaitResult:
问题背景
处理60GB S3源数据,流程为:S3读取→正则过滤无效数据→字段映射→写入Redshift,Job运行2小时后抛出上述报错,原脚本如下:
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 import DynamicFrame import re args = getResolvedOptions(sys.argv, ['JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) AWSGlueDataCatalog_node = glueContext.create_dynamic_frame.from_catalog(database="database_name", table_name="table_name", transformation_ctx="AWSGlueDataCatalog_node") regex_pattern = r'^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}+0000$' Filter_node = Filter.apply(frame = AWSGlueDataCatalog_node, f = lambda x: bool(re.match(regex_pattern, x["col0"])), transformation_ctx="Filter_node") ChangeSchema_node = ApplyMapping.apply(frame=Filter_node, mappings=[("col0", "string", "active_at", "string"), ("col1", "string", "client_id", "string"), ("col2", "string", "profile_uuid", "string")], transformation_ctx="ChangeSchema_node") AmazonRedshift_node = glueContext.write_dynamic_frame.from_options(frame=ChangeSchema_node, connection_type="redshift", connection_options={"redshiftTmpDir": "s3://aws-glue-assets/temporary", "useConnectionProperties": "true", "dbtable": "dbtable", "connectionName": "connectionName", "preactions": "CREATE TABLE IF NOT EXISTS table_name (active_at VARCHAR, client_id VARCHAR, profile_uuid VARCHAR);"}, transformation_ctx="AmazonRedshift_node") job.commit()
可能原因及解决方法
1. 资源不足或写入超时
- 调整Glue Job DPU:默认DPU无法高效处理60GB数据,将DPU数量提升至10-20(根据实际情况调整),增强并行处理能力。
- 优化Spark参数:在Job配置的「Spark参数」中添加:
--conf spark.sql.shuffle.partitions=200:匹配数据量调整分区数,避免分区过多导致任务碎片化,或过少导致单任务负载过高。--conf spark.executor.memory=8g:提升executor内存,减少OOM风险。
- 检查临时目录配置:确保
redshiftTmpDir对应的S3桶与Redshift集群处于同一AWS区域,避免跨区域传输延迟;同时验证Glue执行角色对该S3桶拥有读写权限。
2. 正则过滤效率低下
原脚本使用Python re.match在lambda中做过滤,Python UDF性能远低于Spark内置函数,建议改为Spark SQL正则函数处理:
- 将DynamicFrame转为DataFrame,用
rlike函数过滤,再转回DynamicFrame,大幅提升过滤速度(示例见优化后脚本)。 - 修正正则表达式:原模式中的
+0000应改为\+0000,否则+会被视为量词而非字面字符,导致匹配逻辑错误。
3. Redshift预操作(preactions)耗时过长
preactions中的建表语句会在写入前执行,若表已存在仍会触发检查,可能导致锁表或额外耗时:
- 提前在Redshift手动执行
CREATE TABLE IF NOT EXISTS语句,移除脚本中的preactions参数,减少写入前的额外操作。
4. 数据倾斜问题
若过滤后的数据存在严重倾斜(如某类client_id对应数据量极大),会导致部分任务节点负载过高:
- 通过
df.groupBy("client_id").count().orderBy("count", ascending=False)查看数据分布。 - 对倾斜字段进行加盐处理(如在
client_id后添加随机后缀),分散负载后再写入Redshift。
5. 查看详细错误日志
上述报错仅为上层提示,需查看Glue Job的CloudWatch日志,找到awaitResult对应的具体异常(如Redshift连接超时、OOM、权限错误等),精准定位问题。
优化后的脚本示例
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 import DynamicFrame args = getResolvedOptions(sys.argv, ['JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) # 从Data Catalog读取数据 AWSGlueDataCatalog_node = glueContext.create_dynamic_frame.from_catalog( database="database_name", table_name="table_name", transformation_ctx="AWSGlueDataCatalog_node" ) # 转为DataFrame用Spark内置正则函数过滤,提升性能 df = AWSGlueDataCatalog_node.toDF() # 修正正则表达式,转义+号 regex_pattern = r'^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}\+0000$' filtered_df = df.filter(df.col0.rlike(regex_pattern)) # 转回DynamicFrame继续后续处理 Filter_node = DynamicFrame.fromDF(filtered_df, glueContext, "Filter_node") # 字段映射 ChangeSchema_node = ApplyMapping.apply( frame=Filter_node, mappings=[ ("col0", "string", "active_at", "string"), ("col1", "string", "client_id", "string"), ("col2", "string", "profile_uuid", "string") ], transformation_ctx="ChangeSchema_node" ) # 写入Redshift(提前在Redshift建好表,移除preactions) AmazonRedshift_node = glueContext.write_dynamic_frame.from_options( frame=ChangeSchema_node, connection_type="redshift", connection_options={ "redshiftTmpDir": "s3://aws-glue-assets/temporary", # 确认与Redshift同区域 "useConnectionProperties": "true", "dbtable": "dbtable", "connectionName": "connectionName" }, transformation_ctx="AmazonRedshift_node" ) job.commit()
内容的提问来源于stack exchange,提问作者Nikhil Khandelwal
相关产品推荐
相关产品推荐

