Apache Beam中如何正确创建用于API调用的Session
问题描述
需求为对PCollection中的每个元素发起HTTP API调用。初始实现未使用Session,采用Dataflow runner运行时,处理1200行数据的作业总耗时约9分30秒,其中API调用累计耗时320秒。
为优化API调用性能,尝试在ParDo的setup方法中创建requests Session复用连接,实现代码如下:
class textapi_call(beam.DoFn): def __init__(self, api_key): self.api_key = api_key def setup(self): self.session = requests.session() def process(self, element): address = element[3] + ", " + element[4] + ", " + element[5] + ", " + element[6] + ", " + element[7] url = "https://maps.googleapis.com/maps/api/place/textsearch/json?query=" url += address url += "&key={}".format(api_key) params = {} start = time.time() res = self.session.get(url, params=params) results = json.loads(res.content) time_taken = time.time() - start return [[element[0], address, str(results), time_taken]]
异常现象
- 引入Session后,Dataflow runner环境下作业总耗时仍超过9分钟,API调用累计耗时仍约320秒,无性能提升
- 相同代码在direct runner下运行时,相比未使用Session的版本性能提升超2倍
核心疑问
上述Apache Beam中创建Session的方式是否正确?怀疑当前实现中工作节点上的Session未被正常维持复用。
测试输入示例
AGENT_ID,AGENT_NAME,DATE_OF_JOINING,ADDRESS_LINE1,ADDRESS_LINE2,CITY,STATE,POSTAL_CODE,EMP_ROUTING_NUMBER,EMP_ACCT_NUMBER AGENT00001,Ray Johns,1993-06-05,1402 Maggies Way,,Waterbury Center,VT,05677,034584958,HKUN51252328472585
解答
你在setup方法中初始化Session的写法本身符合Beam的DoFn生命周期规范,但在Dataflow runner下没有拿到预期性能收益,主要有两个原因:
- DoFn实例复用率远低于direct runner环境
setup方法确实是每个DoFn实例初始化时仅执行一次,但direct runner是单进程单实例串行处理全部数据,Session可以全程复用所有连接,所以能看到2倍以上的性能提升。- Dataflow runner会根据负载动态拆分任务、调度worker,默认配置下每个worker、每个DoFn实例实际处理的元素量非常少,很多实例处理几个请求就会被销毁,连接复用率极低,最终累计耗时和每次新建连接没有明显差异。
- 另外你的代码存在一个笔误:拼接URL时使用的是全局变量
api_key而非实例属性self.api_key,该问题不影响Session复用,但容易触发API鉴权错误。
- requests默认Session配置不适合分布式批量调用场景
就算Session被正常复用,requests默认挂载的HTTPAdapter连接池大小仅为10,Dataflow worker默认采用多线程处理元素,并发请求数超过连接池阈值时,依然会新建TCP连接,无法发挥连接复用的优势。
修正实现
你可以按如下方式调整代码,同时配合Dataflow任务参数调整,就能拿到和direct runner一致的性能提升:
import json import time import requests from requests.adapters import HTTPAdapter import apache_beam as beam class textapi_call(beam.DoFn): def __init__(self, api_key): self.api_key = api_key self.session = None self.base_url = None def setup(self): self.session = requests.Session() # 配置连接池参数,pool_maxsize匹配单DoFn实例的最大并发数即可 http_adapter = HTTPAdapter( pool_connections=1, # 仅调用谷歌地图单个域名,连接主机数设为1 pool_maxsize=32, pool_block=True ) self.session.mount("https://", http_adapter) self.base_url = "https://maps.googleapis.com/maps/api/place/textsearch/json" def process(self, element): address = f"{element[3]}, {element[4]}, {element[5]}, {element[6]}, {element[7]}" params = { "query": address, "key": self.api_key } start = time.time() res = self.session.get(self.base_url, params=params, timeout=10) res.raise_for_status() results = res.json() time_taken = time.time() - start return [(element[0], address, str(results), time_taken)] def teardown(self): # DoFn销毁时显式关闭Session释放连接 if self.session: self.session.close()
额外优化建议
- 调整Dataflow任务参数:你的总数据量仅1200条,API累计耗时320秒,串行QPS约为3.75,不需要启动过多worker,限制任务使用单worker、调低单worker并发数,就能让Session连接复用率拉满,总耗时会直接降到和direct runner相当的水平。
- 可以给Session配置重试逻辑,避免偶发网络波动、API限流导致作业失败。
- 谷歌地图Text Search API存在官方QPS限制,不要盲目调高并发触发限流。
内容的提问来源于stack exchange,提问作者Akhil Kv
相关产品推荐
相关产品推荐

