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

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下没有拿到预期性能收益,主要有两个原因:

  1. DoFn实例复用率远低于direct runner环境
    • setup方法确实是每个DoFn实例初始化时仅执行一次,但direct runner是单进程单实例串行处理全部数据,Session可以全程复用所有连接,所以能看到2倍以上的性能提升。
    • Dataflow runner会根据负载动态拆分任务、调度worker,默认配置下每个worker、每个DoFn实例实际处理的元素量非常少,很多实例处理几个请求就会被销毁,连接复用率极低,最终累计耗时和每次新建连接没有明显差异。
    • 另外你的代码存在一个笔误:拼接URL时使用的是全局变量api_key而非实例属性self.api_key,该问题不影响Session复用,但容易触发API鉴权错误。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 04:06:08