AWS Glue作业中如何将线程分发到多个Worker节点运行
问题根源
你当前使用Python原生threading.Thread实现的多线程逻辑,全部运行在Glue作业的单个Driver(主驱动)节点上,完全不会调度到集群内其他Worker节点,这就是你新增多少Worker都无法解决内存问题的核心原因:
- Python原生线程只能在当前代码所在的单个进程内调度,属于单机并行能力,天然感知不到集群中的其他服务器节点
- 你在循环中提前为所有
data_chunk生成json_data的逻辑,本身就会把全量数据都加载到Driver节点的内存中,进一步放大单点内存压力 - 新增的Worker节点全程处于空闲状态,没有被分配任何计算任务,自然无法分担内存压力
解决方案
放弃单机原生多线程的实现,直接使用Glue底层依赖的Spark分布式计算框架原生能力,将数据分片自动分发到所有Worker节点并行处理,实现内存和计算资源的横向扩展。
具体实现代码
from pyspark.sql import SparkSession from awsglue.context import GlueContext import sys from awsglue.utils import getResolvedOptions # 作业基础初始化(Glue作业默认自带该段初始化逻辑,无需重复编写) args = getResolvedOptions(sys.argv, ['JOB_NAME']) spark = SparkSession.builder.getOrCreate() glueContext = GlueContext(spark.sparkContext) # 1. 将本地单机的data_chunks列表转换为分布式RDD,自动切片分发到所有Worker节点 # 可通过numSlices参数指定分片数,建议设置为集群总CPU核数的2~3倍,避免单分片过大撑爆Worker内存 chunks_rdd = spark.sparkContext.parallelize(data_chunks, numSlices=100) # 2. 定义单个数据块的处理逻辑,该函数会被序列化后发送到不同Worker节点执行 def process_single_chunk(data_chunk): # 将json生成逻辑下沉到Worker节点执行,避免全量json堆在Driver内存 json_data = get_bulk_upload_json(data_chunk) # 填入你原来my_func中的业务处理逻辑 process_result = my_func(arg1, arg2, json_data) return process_result # 3. 分布式调度执行:Spark会自动将分片分配到空闲Worker上运行,内存压力分散到所有节点 # 如果不需要把全量结果汇总回Driver,不要调用collect(),直接将结果写入S3/目标库即可 results = chunks_rdd.map(process_single_chunk).collect()
关键注意事项
- 禁止在Driver端提前做全量数据预处理:所有针对单个数据块的处理、转换逻辑,都要下沉到
process_single_chunk这类分发到Worker执行的函数中,避免全量数据堆积在Driver节点。 - 合理设置分片大小:如果单条
data_chunk本身数据量很大,要适当调大numSlices参数,保证每个分片的内存占用不超过单个Worker可分配的内存阈值。 - 单Worker内部的线程优化仅作为补充:如果你的业务逻辑是IO密集型(比如调用第三方接口、上传文件到S3),可以在
process_single_chunk函数内部开启少量线程提升单Worker的IO利用率,但这一步的前提是任务已经通过Spark调度到不同Worker上,否则依然会出现单点内存瓶颈。 - 避免全量结果拉回Driver:如果处理结果不需要在Driver端做全局汇总,不要调用
collect()方法,直接在分布式计算过程中将结果写入目标存储(比如S3、JDBC数据库、Glue数据表),避免结果集过大撑爆Driver内存。
内容的提问来源于stack exchange,提问作者rodrigocf
相关产品推荐
相关产品推荐

