如何在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
相关产品推荐
相关产品推荐

