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)
问题原因分析
- 数据分区不均匀:哈希分区的字段选择不当,导致部分Executor负载过高,内存/CPU耗尽,无法维持心跳连接
- 资源配置不足:Glue Worker的内存或CPU配额不足以处理百万级数据的转换计算,引发超时
- 代码拼写错误:转换代码中
tragetdf的拼写错误可能导致变量混乱,大数据量下放大资源消耗问题 - 过早转换与分区:读取后立即转换为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
相关产品推荐
相关产品推荐

