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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 08:07:41