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

Spark超时错误:Glue DynamicFrame转Spark DataFrame失败

问题:Glue任务读取百万级Aurora数据时转换超时失败

场景重现

使用Glue任务从Aurora DB读取数据,小量数据(约5万条)可正常运行,但数据量达到数百万条时,在将DynamicFrame转换为Spark DataFrame的步骤触发超时错误。

读取数据代码

df = self.glue_context.create_dynamic_frame.from_options(
            connection_type="custom.jdbc",
            connection_options={
                "className": self.jdbc_driver_name,
                "url": self.aurora_url,
                "user": self.db_username,
                "password": self.db_password,
                "query": query,
                "hashexpression": hash_expression,
                "hashpartitions": hash_partition,
            },
        )

转换与持久化代码

targetdf = df.toDF()
# 注意:存在拼写错误,tragetdf应为targetdf
tragetdf = tragetdf.select(col("col1").alias("col1"),
                           col("col2").alias("col2")
                          ).repartition(int(partitions))
tragetdf.persist()

错误日志

Traceback (most recent call last):
  File "/tmp/code.py", line 857, in <module>
    init()
  File "/tmp/code.py", line 853, in init
    main(args, glue_context, spark, current_path, bucket_name, file_path, fixed_path, i_output, final_output_path)
  File "/tmp/code.py", line 362, in main
    brk_df = fetch_from_aurora(args, glue_context, source_table, hash_expression_aurora, target_query)
  File "/tmp/code.py", line 274, in fetch_from_aurora
    df = intermtntdf.toDF()
  File "/opt/amazon/lib/python3.6/site-packages/awsglue/dynamicframe.py", line 148, in toDF
    return DataFrame(self._jdf.toDF(self.glue_ctx._jvm.PythonUtils.toSeq(scala_options)), self.glue_ctx)
  File "/opt/amazon/spark/python/lib/py4j-0.10.7-src.zip/py4j/java_gateway.py", line 1257, in __call__
    answer, self.gateway_client, self.target_id, self.name)
  File "/opt/amazon/spark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 63, in deco
    return f(*a, **kw)
  File "/opt/amazon/spark/python/lib/py4j-0.10.7-src.zip/py4j/protocol.py", line 328, in get_return_value
    format(target_id, ".", name), value)
py4j.protocol.Py4JJavaError: An error occurred while calling o135.toDF.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 8.0 failed 4 times, 
most recent failure: Lost task 0.3 in stage 8.0 (TID 11, 10.156.19.74, executor 12):
 ExecutorLostFailure (executor 12 exited caused by one of the running tasks) Reason: Executor heartbeat timed out after 648979 ms
Driver stacktrace:
	at org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1889)
	at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1877)
	at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1876)
	at scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
	at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48)

问题原因分析

  1. 数据分区不均匀:哈希分区的字段选择不当,导致部分Executor负载过高,内存/CPU耗尽,无法维持心跳连接
  2. 资源配置不足:Glue Worker的内存或CPU配额不足以处理百万级数据的转换计算,引发超时
  3. 代码拼写错误:转换代码中tragetdf的拼写错误可能导致变量混乱,大数据量下放大资源消耗问题
  4. 过早转换与分区:读取后立即转换为DataFrame并repartition,大量数据在内存中集中处理,引发资源瓶颈

优化方案

1. 优化读取分区策略

  • 确保hash_expression选择分布均匀的字段(如主键、唯一键),避免数据倾斜
  • 调整hashpartitions数量,建议与Glue Worker数量匹配(每个Worker对应2-4个分区)
  • 尝试改用范围分区替代哈希分区,适合数值型分布均匀的字段:
    df = self.glue_context.create_dynamic_frame.from_options(
        connection_type="custom.jdbc",
        connection_options={
            "className": self.jdbc_driver_name,
            "url": self.aurora_url,
            "user": self.db_username,
            "password": self.db_password,
            "dbtable": source_table,
            "partitionColumn": "id",  # 替换为范围分布均匀的字段
            "lowerBound": "1",
            "upperBound": "1000000",
            "numPartitions": "20"
        },
    )
    

2. 调整Glue资源配置

  • 升级Worker实例类型(如从标准型改为G.2X),提升单Worker的内存与CPU
  • 增加Worker数量,分散数据处理压力
  • 添加Spark参数优化资源配额:
    --conf spark.executor.memory=16g --conf spark.executor.cores=4 --conf spark.driver.memory=8g
    

3. 修正代码逻辑

  • 统一修正变量拼写错误,将tragetdf改为targetdf
  • 先通过Glue DynamicFrame做数据裁剪(如SelectFields),减少转换后的数据量:
    trimmed_df = df.select_fields(["col1", "col2"])
    targetdf = trimmed_df.toDF().repartition(int(partitions))
    
  • 按需使用persist():若后续无需多次复用DataFrame,可移除该操作;若需要,指定更高效的存储级别(如MEMORY_AND_DISK_SER)
  • 直接使用Spark JDBC读取,跳过DynamicFrame层,更灵活控制读取逻辑:
    targetdf = spark.read.jdbc(
        url=self.aurora_url,
        table=f"({query}) as tmp_table",
        properties={
            "user": self.db_username,
            "password": self.db_password,
            "driver": self.jdbc_driver_name
        },
        partitionColumn=hash_expression,
        numPartitions=int(hash_partition)
    )
    

4. 排查数据倾斜

  • 通过Spark UI查看Stage任务的处理数据量,定位是否存在单个Task处理远超平均数据量的情况
  • 对倾斜字段做加盐处理,或拆分查询分批读取数据

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 12:05:38