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

AWS Glue 4.0任务中Python多进程执行挂起问题求助

问题描述

我尝试在同一个AWS Glue 4.0任务中使用Python Multiprocessing并行处理数据,因特定原因不想采用Glue Workflows多任务方式实现并行。

我的Python代码

from multiprocessing import Pool
import sys
import time
import random

from awsglue.context import GlueContext
from awsglue.job import Job
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext

args = getResolvedOptions(sys.argv, ['JOB_NAME', 'TempDir'])

sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)
print(f"{args['JOB_NAME']} STARTED")

def worker(table_name, tmp_dir):
    print(f"STARTED WORKER: {table_name}")
    data = load_data(table_name, tmp_dir)
    process_data(table_name, data)
    print(f"FINISHED WORKER: {table_name}")
    
def load_data(table_name, tmp_dir):    
    print(f"LOADING: {table_name}")
    data = glueContext.create_dynamic_frame.from_catalog(database="my_database",
                                                         table_name=table_name,
                                                         redshift_tmp_dir=f"{tmp_dir}/{table_name}",
                                                         transformation_ctx=f"data_source_{table_name}")
    time.sleep(random.randint(1, 5))  # added here to simulate different loading times
    print(f"LOADED: {table_name} has {data.count()} rows")
    return data

def process_data(table_name, data):
    print(f"PROCESSING: {table_name}")
    # do something
    time.sleep(random.randint(1, 5))  # added here to simulate different processing times
    print(f"DONE: {table_name}")

pool = Pool(4)
tables = ['TABLE1', 'TABLE2', 'TABLE3', 'TABLE4', 'TABLE5', 'TABLE6', 'TABLE7', 'TABLE8', 'TABLE9']
for table in tables:
    pool.apply_async(worker, args=(table, args['TempDir']))
pool.close()
pool.join()

print(f"{args['JOB_NAME']} COMPLETED")
job.commit()

问题现象

任务看似正常启动多个worker后就挂起,直到超时才结束。CloudWatch输出日志如下,无错误日志:

2023-04-19T12:01:49.566+02:00   STARTED WORKER: TABLE1 LOADING: TABLE1
2023-04-19T12:01:49.566+02:00   STARTED WORKER: TABLE2 LOADING: TABLE2
2023-04-19T12:01:49.566+02:00   STARTED WORKER: TABLE3 LOADING: TABLE3 
2023-04-19T12:01:49.566+02:00   STARTED WORKER: TABLE4 LOADING: TABLE4
2023-04-19T12:01:49.603+02:00   STARTED WORKER: TABLE5 LOADING: TABLE5
2023-04-19T12:01:49.604+02:00   STARTED WORKER: TABLE6 LOADING: TABLE6
2023-04-19T12:01:49.607+02:00   STARTED WORKER: TABLE7 LOADING: TABLE7
2023-04-19T12:01:49.608+02:00   STARTED WORKER: TABLE8 LOADING: TABLE8
2023-04-19T12:01:49.609+02:00   STARTED WORKER: TABLE9 LOADING: TABLE9

我尝试了多种方法,但仅能确定问题似乎出在create_dynamic_frame.from_catalog()环节。请问为何无法正常运行?有人遇到并解决过该问题吗?


解决方案与原因分析

核心原因

AWS Glue的GlueContext、SparkContext这类Spark相关的上下文对象不支持跨Python进程共享。当你用multiprocessing.Pool创建子进程时,这些上下文对象会被尝试序列化/复制到子进程中,但Spark的上下文设计是单进程单实例的,子进程中的上下文处于无效状态,调用create_dynamic_frame.from_catalog()时会陷入无限等待(因为无法和Spark集群建立有效通信)。

另外,Python的multiprocessing在fork进程时,Spark的底层资源连接(比如JVM通信管道)会被破坏,导致子进程无法正常执行Spark相关操作。

替代实现方案

方案1:利用Spark的并行化特性处理多个表

将表名列表转为RDD/DataFrame,用Spark的分布式执行能力并行处理每个表,这是Glue任务内并行处理的最优解:

from awsglue.context import GlueContext
from awsglue.job import Job
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
import sys
import time
import random

args = getResolvedOptions(sys.argv, ['JOB_NAME', 'TempDir'])

sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)
print(f"{args['JOB_NAME']} STARTED")

def process_table(table_name):
    print(f"STARTED PROCESSING: {table_name}")
    # 加载数据
    print(f"LOADING: {table_name}")
    data = glueContext.create_dynamic_frame.from_catalog(
        database="my_database",
        table_name=table_name,
        redshift_tmp_dir=f"{args['TempDir']}/{table_name}",
        transformation_ctx=f"data_source_{table_name}"
    )
    time.sleep(random.randint(1, 5))
    print(f"LOADED: {table_name} has {data.count()} rows")
    
    # 处理数据
    print(f"PROCESSING: {table_name}")
    time.sleep(random.randint(1, 5))
    print(f"DONE: {table_name}")
    return f"Success:{table_name}"

# 将表名转为并行RDD,设置并行度控制并发数
tables = ['TABLE1', 'TABLE2', 'TABLE3', 'TABLE4', 'TABLE5', 'TABLE6', 'TABLE7', 'TABLE8', 'TABLE9']
parallel_rdd = sc.parallelize(tables, numSlices=4)
# 用map执行并行处理
results = parallel_rdd.map(process_table).collect()

print(f"PROCESS RESULTS: {results}")
print(f"{args['JOB_NAME']} COMPLETED")
job.commit()

方案2:使用Python多线程处理(适合IO密集型场景)

如果不需要分布式并行,只是想在Glue驱动进程内用多线程处理,可替换multiprocessing为threading(Spark上下文对象是线程安全的):

from threading import Thread
import sys
import time
import random
from awsglue.context import GlueContext
from awsglue.job import Job
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext

args = getResolvedOptions(sys.argv, ['JOB_NAME', 'TempDir'])

sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)
print(f"{args['JOB_NAME']} STARTED")

def worker(table_name, tmp_dir):
    print(f"STARTED WORKER: {table_name}")
    print(f"LOADING: {table_name}")
    data = glueContext.create_dynamic_frame.from_catalog(
        database="my_database",
        table_name=table_name,
        redshift_tmp_dir=f"{tmp_dir}/{table_name}",
        transformation_ctx=f"data_source_{table_name}"
    )
    time.sleep(random.randint(1, 5))
    print(f"LOADED: {table_name} has {data.count()} rows")
    
    print(f"PROCESSING: {table_name}")
    time.sleep(random.randint(1, 5))
    print(f"DONE: {table_name}")
    print(f"FINISHED WORKER: {table_name}")

tables = ['TABLE1', 'TABLE2', 'TABLE3', 'TABLE4', 'TABLE5', 'TABLE6', 'TABLE7', 'TABLE8', 'TABLE9']
threads = []
# 控制并发线程数
max_concurrent = 4
for i in range(0, len(tables), max_concurrent):
    batch = tables[i:i+max_concurrent]
    for table in batch:
        t = Thread(target=worker, args=(table, args['TempDir']))
        threads.append(t)
        t.start()
    # 等待当前批次线程完成
    for t in threads[-max_concurrent:]:
        t.join()

print(f"{args['JOB_NAME']} COMPLETED")
job.commit()

关键注意事项

  • 永远不要在Glue任务中用Python多进程处理Spark/Glue相关操作,这违背了Spark的架构设计。
  • 用Spark并行化的方案适合大数据量处理,能利用Glue集群的分布式资源。
  • 多线程方案需控制并发数,避免触发Glue或数据源的限流。

内容的提问来源于stack exchange,提问作者Roberto A.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 10:55:22