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

如何在DataBricks中高效并行处理12M条API请求以缩短耗时

优化Databricks中大规模单条API数据拉取至日级处理的方案

问题背景

当前项目需处理包含1200多万条ID的DataFrame(df_ids),每条ID必须通过仅支持单条查询的API拉取数据,处理函数如下:

def ELOQUA_CONTACT(id):
  API = EloquaAPI(f'1.0/data/contact/{id}')
  try:
    contactid = API['id'].lower()
  except:
    contactid = ''
  try:
    company = API['accountName']
  except:
    company = ''

  df = pd.DataFrame([contactid, company]).T.rename(columns={0:'contactid', 1:'company'}) 
  return df

已在AWS Databricks环境测试两种Python并行方案,但效率远未达标:

  • 自定义线程类方案:处理1000条耗时近2分钟,推算处理1200万条需约916天;
  • ThreadPool方案:处理1000条耗时约20秒,推算处理1200万条需约167天;

需求:基于Databricks、AWS、Spark技术栈优化,实现日级处理,最终部署为Databricks调度任务,可调整集群CPU/内存资源。


优化方案

一、Spark分布式并行优化(核心方向)

1. 封装Spark UDF实现分布式执行

利用Spark的集群分布式能力,将任务拆分到多个Executor节点并行处理:

  • 先对df_ids按Executor数量合理分区(建议每个分区1000-5000条,根据API限流调整);
  • 定义PySpark UDF并处理异常与超时:
from pyspark.sql.functions import udf
from pyspark.sql.types import StructType, StructField, StringType

# 定义返回Schema
schema = StructType([
    StructField("contactid", StringType(), True),
    StructField("company", StringType(), True)
])

@udf(returnType=schema)
def get_eloqua_contact(id):
    try:
        API = EloquaAPI(f'1.0/data/contact/{id}')
        contactid = API['id'].lower() if 'id' in API else ''
        company = API['accountName'] if 'accountName' in API else ''
        return (contactid, company)
    except Exception:
        return ('', '')

# 执行UDF并拆分结果
result_df = df_ids.withColumn("contact_data", get_eloqua_contact("id")) \
                  .select("id", "contact_data.contactid", "contact_data.company")
  • 调整Spark配置:增加Executor数量(如spark.executor.instances=30)、每个Executor核心数(如spark.executor.cores=4),同时设置API调用超时时间,避免单个请求阻塞任务。

2. 批量异步调用+重试机制

在每个Executor节点内用异步请求提升单节点处理效率,同时添加重试应对API临时失败:

import aiohttp
import asyncio
from tenacity import retry, stop_after_attempt, wait_exponential
from pyspark.sql.functions import pandas_udf
import pandas as pd

async def fetch_contact(session, id):
    @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10))
    async def _fetch():
        async with session.get(f"你的API基础地址/1.0/data/contact/{id}") as response:
            if response.status == 200:
                data = await response.json()
                contactid = data['id'].lower() if 'id' in data else ''
                company = data['accountName'] if 'accountName' in data else ''
                return (id, contactid, company)
            else:
                raise Exception(f"API请求失败,状态码:{response.status}")
    try:
        return await _fetch()
    except:
        return (id, '', '')

async def batch_fetch(ids):
    async with aiohttp.ClientSession() as session:
        tasks = [fetch_contact(session, id) for id in ids]
        results = await asyncio.gather(*tasks)
        return results

# 封装为Pandas UDF处理批量数据
@pandas_udf("id string, contactid string, company string")
def batch_get_contacts(ids: pd.Series) -> pd.DataFrame:
    loop = asyncio.get_event_loop()
    results = loop.run_until_complete(batch_fetch(ids.tolist()))
    return pd.DataFrame(results, columns=["id", "contactid", "company"])

# 按分区批量处理
result_df = df_ids.groupBy(spark_partition_id()).applyInPandas(batch_get_contacts, schema="id string, contactid string, company string")

注:需在Databricks集群安装aiohttp、tenacity依赖

二、AWS与Databricks集群资源优化

  • 选用计算优化型实例:如c5.4xlarge,这类实例CPU性能强,适合高并发API调用场景;
  • 开启自动缩放:配置Databricks集群自动缩放规则,根据任务负载动态增减Executor数量,平衡资源利用率与处理效率;
  • 缓存输入数据:执行df_ids.cache()将输入DataFrame缓存至内存,避免重复读取开销;
  • 配置Job调度策略:将任务部署为Databricks Job,设置任务重试次数、资源池优先级,确保任务稳定运行。

三、API调用细节优化

  • 复用连接池:在HTTP请求中启用连接池,避免重复建立TCP连接,减少网络开销;
  • 严格控制并发数:根据API的Rate Limit调整每个Executor的异步任务数,避免被限流或封禁;
  • 异步替代线程:用异步IO替代多线程,减少线程切换开销,提升单节点处理密度。

内容的提问来源于stack exchange,提问作者Andrew Hicks

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 08:55:18