Azure Synapse中PySpark循环调用API耗时过长问题排查
Azure Synapse PySpark笔记本API调用耗时排查
我在Azure Synapse Analytics的PySpark笔记本中调用限制为每分钟30次的API(未并行处理),但操作耗时极长,而本地PY3内核的Jupyter笔记本中相同代码速度极快。作为PySpark新手,想排查哪些操作导致了耗时问题。
请求函数代码
# Request functions # First function def get_response(url): headers = {'Authorization': f'Bearer {INSEE_KEY}'} response = requests.get(url, headers=headers) return response #Second function using the first one def get_data_by_siret(siret): siret_url = f'{base_url}siret?q=siret:{siret}' response = get_response(siret_url) if response.status_code == 200 : content = response.json() siret_data = content['etablissements'] print(f'Siret: {siret} collecté. (timestamp: {time.time()})') return siret_data, 200 elif response.status_code == 429: print(f'Siret: {siret} ; quota atteint. (timestamp: {time.time()})') return [], 429 elif response.status_code == 404: print(f'Siret {siret} inconnu dans la base Sirene') return [], 404 else : print(f'Siret {siret}. {response.status_code}, (timestamp: {time.time()})') return [], response.status_code
耗时严重的循环代码
# Collecting Sirets data sirets_json = [] timestamps = [] siret_429 = [] for siret in array_siret: time.sleep(1) siret_data, status = get_data_by_siret(siret) timestamps.append(time.time()) if status == 200: sirets_json.extend(siret_data) elif status == 429: time.sleep(1) siret_429.append(siret)
耗时表现示例
Siret: 04555025800029 collected. (timestamp: 1720795818.7682683) 'Fri Jul 12 14:50:18 2024'
Siret: 05780615000017 collected. (timestamp: 1720795851.364306) 'Fri Jul 12 14:50:51 2024'
Siret: 09572031400475 collected. (timestamp: 1720795883.9143667) 'Fri Jul 12 14:51:23 2024'
本地环境中,相同代码可在约75秒内完成60条Siret数据的采集。
排查方向
- 网络延迟差异:Azure Synapse集群节点与API服务器的网络链路可能比本地更长,存在更高延迟。可在
get_response函数中添加单次请求耗时统计,对比本地和Synapse环境的响应时间。 - 驱动节点资源限制:PySpark笔记本的驱动节点CPU、内存不足,导致循环执行时出现资源争抢或上下文切换开销。查看Synapse工作区的驱动节点配置,对比本地机器资源情况。
print语句的额外开销:PySpark环境中,驱动节点的print输出需要经过集群日志系统转发,频繁打印会增加不必要的开销。暂时注释所有print语句,测试耗时是否下降。requests库环境差异:Synapse环境的requests库版本或依赖(如SSL库)与本地不同,导致请求处理效率降低。对比两边的库版本,或尝试改用PySpark原生HTTP请求方式(若API支持批量)。- Python-JVM交互开销:PySpark驱动节点执行Python循环时,存在JVM与Python进程的交互开销。可尝试将循环逻辑改为PySpark的RDD/DataFrame操作(如
map函数),同时注意控制并发以符合API限流要求。 time.sleep精度问题:Synapse环境中time.sleep的实际等待时长可能受集群调度影响,超过预期值。在每次sleep前后记录时间,验证实际等待时长。
内容的提问来源于stack exchange,提问作者Greg Mansio
相关产品推荐
相关产品推荐

