AWS EMR运行PySpark任务调用mapPartitions时卡住如何排查?
问题原因分析
1. 最可能的根因:集群网络权限问题
本地环境能访问目标接口,但AWS集群的Worker节点默认可能没有公网出口权限:
- Worker节点所在的私有子网未配置NAT网关,无法访问公网接口
- 安全组出站规则未放开目标接口的80/443端口访问权限
- 你代码中
requests.post未设置超时时间,网络不通时会无限等待,表现为任务卡住
2. 代码语法/逻辑缺陷
你提供的代码存在多处语法和逻辑问题,本地能运行大概率是粘贴时的笔误,集群运行时会触发隐性错误:
some_func1拼写错误(代码中写为some_fuc1)- 类方法缺少
self参数,send_them、extr_data调用时会出现参数不匹配问题 requests.post语法错误,多写了右括号,且headers未按字典格式传入、请求URL缺少http/https前缀- 未捕获请求异常,出现连接错误时直接挂住不会抛出日志
3. 闭包序列化问题
你传入mapPartitions的lambda直接引用了驱动端的dict_name和id变量,极端情况下会因为闭包序列化失败导致任务卡住,不过该问题概率低于网络问题。
调试步骤
- 优先验证Worker节点网络连通性:登录集群Worker节点,执行
curl -v 你的目标接口地址,确认是否能正常访问 - 给
requests.post加上超时参数timeout=10,同时添加异常捕获和日志打印,卡住的任务会直接抛出具体错误信息,可到Spark UI的Executor日志页查看报错 - 临时修改
some_func1为直接返回固定值,跳过接口请求:如果任务能正常运行,即可100%确定是接口请求相关的问题,和Spark逻辑无关 - 小数据量场景下可以先把RDD结果
collect到驱动端,本地循环调用some_func1验证逻辑是否正常
修复方案
1. 优先修复代码语法和请求逻辑
修改some_func1添加超时和异常捕获:
import requests import json def some_func1(rec, dict_name, id): try: rec_list = list(rec) # 注意headers要按字典格式传入,不能是字符串 headers = {"Content-Type": "application/json"} # 必须加http/https前缀 attrburl = "https://www.someurl.com" response = requests.post( attrburl, data=json.dumps(rec_list), headers=headers, # 10秒连接+读取超时,避免无限卡住 timeout=10 ) # 遇到4xx/5xx错误直接抛出异常 response.raise_for_status() return response.json() except Exception as e: # 打印错误信息到Executor日志,方便排查 print(f"请求失败,错误信息:{str(e)},请求数据:{rec_list}") return {"error": str(e)}
2. 修复类方法逻辑
补全self参数,修正方法调用:
class Processor: def __init__(self, sc, arguments): self.sc = sc self.env = arguments.env self.dte = arguments.dte self.sendme = arguments.sendme def send_them(self, ext_data, dict_name, id): attributes = ext_data.rdd.map(lambda x: ctgs(x['col1'], x['col2'], x['col3'])) # 闭包问题修复:改用广播变量传递参数,避免序列化问题 dict_bc = self.sc.sparkContext.broadcast(dict_name) id_bc = self.sc.sparkContext.broadcast(id) def process_partition(iter): return [some_func1(map(lambda x: x, xs), dict_bc.value, id_bc.value) for xs in partition_all(50, iter)] response = attributes.mapPartitions(process_partition).collect() # 销毁广播变量释放资源 dict_bc.destroy() id_bc.destroy() return response def extr_data(self, dict_name, id): ext_data = self.sc.sql('''select col1, col2, col3 from table_name''') return self.send_them(ext_data, dict_name, id) def process(self): dict_name = { "dict_id": '34343-3434-3433-343'} id = 'dfdfd-erere-dfd' self.extr_data(dict_name, id)
3. 小数据量场景简化方案
你只有100条数据,完全可以直接在驱动端发送请求,绕过Worker节点的网络限制:
def send_them(self, ext_data, dict_name, id): # 直接把数据拉取到驱动端,不需要走Worker节点计算 attributes = ext_data.rdd.map(lambda x: ctgs(x['col1'], x['col2'], x['col3'])).collect() response = [] for xs in partition_all(50, attributes): response.append(some_func1(xs, dict_name, id)) return response
4. 集群网络配置修复
如果确认是Worker网络不通,根据接口地址类型配置对应权限:
- 公网接口:给Worker所在子网配置NAT网关,放开安全组出站443端口权限
- AWS内部服务接口:配置对应VPC端点,放开安全组内部访问权限
内容的提问来源于stack exchange,提问作者user7343922
相关产品推荐
相关产品推荐

