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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 07:23:14