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

Cloud Functions (Gen 2) API请求策略是否合理?技术方案咨询

问题背景与架构咨询

我正尝试构建由Cloud Scheduler触发的GCP Cloud Functions (Gen 2),用于从博彩商API采集不同赛程的各类体育赛事赔率数据,采用Python开发,使用GCP的Cloud Scheduler、Pub/Sub、Firestore服务。

目标

从博彩商API获取赔率数据,并将所需数据点存储至Firestore。

计划

为不同赛程的体育赛事配置Cloud Functions (Gen 2)与Cloud Scheduler,实现赔率数据采集。

当前架构设计

每项赛事需采集3-5个投注市场的数据,部分博彩商API需先调用接口获取event ids,再通过这些event ids调用接口获取赛事赔率数据,完成所有投注市场数据采集需3-5次API调用。目前本地将不同API拆分为独立函数可正常运行,现计划采用一对多Pub/Sub消息系统部署这些函数:

  • Function 1:采集event ids并推送至Pub/Sub主题
  • Function 2:订阅Pub/Sub消息,使用event ids采集赔率市场1数据
  • Function 3:订阅Pub/Sub消息,使用event ids采集赔率市场2数据

要求Function 1先由Cloud Scheduler触发,采集数据并推送到Pub/Sub后,再触发Function 2、3执行。

Function 1(赛事ID采集)代码

import requests
import json
from firebase_admin import firestore
from google.cloud import pubsub_v1

db = firestore.Client(project='xxxxxxx')

# API INFO
Base_url = 'https://xxxxxxx.net/v1/feeds/sportsbookv2'
Sport_id = '1000093204'
AppID = 'xxxxxxx'
AppKey = 'xxxxxxx'
Country = 'en_AU'
Site = 'xxxxxxx'


project_id = "xxxxxxx"
topic_id = "xxxxxxx_basketball_nba"


publisher = pubsub_v1.PublisherClient()
# The `topic_path` method creates a fully qualified identifier
# in the form `projects/{project_id}/topics/{topic_id}`
topic_path = publisher.topic_path(project_id, topic_id)


def win_odds(request):

    event_ids = []

    url = f"{Base_url}/betoffer/event/{','.join(map(str, event_ids))}.json?app_id={AppID}&app_key={AppKey}&local={Country}&site={Site}"
    print(url)

    windata = requests.get(url).text
    windata = json.loads(windata)

    for odds_data in windata['betOffers']:
        if odds_data['betOfferType']['name'] == 'Head to Head' and 'MAIN' in odds_data['tags']:
            event_id = odds_data['eventId']
            home_team = odds_data['outcomes'][0]['participant']
            home_team_win_odds = odds_data['outcomes'][0]['odds']
            away_team = odds_data['outcomes'][1]['participant']
            away_team_win_odds = odds_data['outcomes'][1]['odds']

            print(f'{event_id} {home_team} {home_team_win_odds} {away_team} {away_team_win_odds}')

            event_ids.append(event_id)

            # WRITE TO FIRESTORE
            doc_ref = db.collection(u'xxxxxxx_au').document(u'basketball_nba').collection(u'win_odds').document(
                f'{event_id}')
            doc_ref.set({
                u'event_id': event_id,
                u'home_team': home_team,
                u'home_team_win_odds': home_team_win_odds,
                u'away_team': away_team,
                u'away_team_win_odds': away_team_win_odds,
                u'timestamp': firestore.SERVER_TIMESTAMP,
            })

    return event_ids


print(f"Published messages to {topic_path}.")

订阅函数模板代码

from concurrent.futures import TimeoutError
from google.cloud import pubsub_v1
import requests
import json
from firebase_admin import firestore
import functions_framework


db = firestore.Client(project='xxxxxxx')

# API INFO
Base_url = 'https://xxxxxxx.net/v1/feeds/sportsbookv2'
Sport_id = '1000093204'
AppID = 'xxxxxxx'
AppKey = 'xxxxxxx'
Country = 'en_AU'
Site = 'xxxxxxx'

project_id = "xxxxxxx"
subscription_id = "xxxxxxx_basketball_nba_events"
timeout = 5.0

subscriber = pubsub_v1.SubscriberClient()
subscription_path = subscriber.subscription_path(project_id, subscription_id)

@functions_framework.cloud_event
def win_odds(message):
    print(message.data)

    event_ids = json.loads(message.data)

    url = f"{Base_url}/betoffer/event/{','.join(map(str, event_ids))}.json?app_id={AppID}&app_key={AppKey}&local={Country}&site={Site}"
    print(url)

    windata = requests.get(url).text
    windata = json.loads(windata)

    for odds_data in windata['betOffers']:
        if odds_data['betOfferType']['name'] == 'Head to Head' and 'MAIN' in odds_data['tags']:
            event_id = odds_data['eventId']
            home_team = odds_data['outcomes'][0]['participant']
            home_team_win_odds = odds_data['outcomes'][0]['odds']
            away_team = odds_data['outcomes'][1]['participant']
            away_team_win_odds = odds_data['outcomes'][1]['odds']

            print(f'{event_id} {home_team} {home_team_win_odds} {away_team} {away_team_win_odds}')

            # WRITE TO FIRESTORE
            doc_ref = db.collection(u'xxxxxxx_au').document(u'basketball_nba').collection(u'win_odds').document(
                f'{event_id}')
            doc_ref.set({
                u'event_id': event_id,
                u'home_team': home_team,
                u'home_team_win_odds': home_team_win_odds,
                u'away_team': away_team,
                u'away_team_win_odds': away_team_win_odds,
                u'timestamp': firestore.SERVER_TIMESTAMP,
            })


if __name__ == '__main__':
    streaming_pull_future = subscriber.subscribe(subscription_path, callback=win_odds)
    print(f"Listening for messages on {subscription_path}..\n")

    # Wrap subscriber in a 'with' block to automatically call close() when done.
    with subscriber:
        try:
            # When `timeout` is not set, result() will block indefinitely,
            # unless an exception is encountered first.
            streaming_pull_future.result(timeout=timeout)
        except TimeoutError:
            streaming_pull_future.cancel()  # Trigger the shutdown.
            streaming_pull_future.result()  # Block until the shutdown is complete.

    print("Listening, then shutting down.")

请问该策略是否合理,是否存在更优的实现方案?


回答

现有策略的合理性分析

你的架构思路是合理的,核心优势包括:

  1. 解耦性:通过Pub/Sub拆分事件ID采集与各投注市场的赔率采集逻辑,避免单一函数复杂度过高,后续新增投注市场时只需添加新的订阅函数即可,扩展性强。
  2. 异步执行:Pub/Sub支持异步消息传递,Function1完成采集后无需等待后续函数执行,可立即终止,提升整体流程效率。
  3. 故障隔离:单个投注市场的采集失败不会影响其他市场或事件ID的采集流程,降低系统整体故障风险。

但现有代码存在几个关键问题需要修正:

  • Function1未实际推送消息到Pub/Sub:代码仅初始化了PublisherClient,但未调用publisher.publish()方法将event_ids发送到主题,导致后续订阅函数无法收到触发信号。
  • 订阅函数的本地逻辑与云函数部署冲突:云函数模式下,GCP会自动处理Pub/Sub消息的触发,无需手动创建SubscriberClient并启动监听,本地调试的if __name__ == '__main__'块需在部署时移除或注释。
  • Function1初始API调用无效:event_ids初始为空列表,拼接后的URL请求空事件集合,无法获取有效数据,应先调用赛事列表API拿到event_ids,再请求赔率数据。

更优实现方案建议

方案1:优化现有Pub/Sub架构

关键调整点:

  1. 修复Function1的消息推送逻辑:在win_odds函数末尾添加消息推送代码(注意Pub/Sub要求数据为字节类型):
# 在Function1的win_odds函数末尾添加
import json

if event_ids:
    data = json.dumps(event_ids).encode("utf-8")
    future = publisher.publish(topic_path, data)
    print(f"Published message ID: {future.result()}")
  1. 简化订阅函数的云函数逻辑:移除本地订阅监听代码,专注处理Cloud Event触发的消息:
# 订阅函数简化后
import base64
import json
from functions_framework import cloud_event

@cloud_event
def win_odds(cloud_event):
    # 解析Pub/Sub消息数据
    message_data = base64.b64decode(cloud_event.data["message"]["data"]).decode("utf-8")
    event_ids = json.loads(message_data)
    
    # 后续赔率采集与写入Firestore逻辑保持不变
    # ...
  1. 添加错误处理与重试机制:
    • 对API调用添加超时、重试(可使用tenacity库),避免网络波动导致采集失败。
    • 配置Pub/Sub订阅的死信队列,将处理失败的消息转发到死信主题,便于后续排查与重处理。
  2. 批量处理优化:若event_ids数量较多,拆分小批量进行API调用,避免单次请求参数过长被API拒绝。

方案2:使用Cloud Workflows编排流程

如果需要更严格的流程控制(例如确保所有投注市场采集完成后执行后续操作),可以用Cloud Workflows替代Pub/Sub的一对多触发,实现精细的流程编排:

  1. Cloud Scheduler触发Cloud Workflows。
  2. Workflows第一步调用Function1获取event_ids。
  3. Workflows并行调用Function2、Function3传递event_ids。
  4. 可添加流程分支、错误处理、重试逻辑,甚至在所有函数执行完成后触发通知或数据校验任务。

该方案优势是流程可视化、可追溯,适合对执行顺序和结果有明确要求的场景,但配置复杂度略高于Pub/Sub方案。

方案3:单函数多任务并行处理(适合小规模场景)

如果投注市场数量较少(3-5个),且API调用频率限制宽松,可将所有逻辑合并到单个Cloud Function中,使用Python的concurrent.futures.ThreadPoolExecutor并行调用不同投注市场的API:

import concurrent.futures

def fetch_market_odds(event_ids, market_type):
    # 根据market_type调用对应API获取赔率并写入Firestore
    # ...

def main(request):
    # 第一步:获取event_ids
    event_ids = fetch_event_ids()
    
    # 第二步:并行采集各市场数据
    with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor:
        futures = [executor.submit(fetch_market_odds, event_ids, market) for market in ["Head to Head", "Over/Under", ...]]
        for future in concurrent.futures.as_completed(futures):
            future.result()

该方案优势是部署简单、减少云资源调用次数,但函数逻辑相对复杂,故障影响范围更大,适合小规模、低复杂度的采集需求。

总结

  • 现有Pub/Sub架构思路合理,修正代码中的推送逻辑和云函数适配问题后即可投入使用,适合需要解耦、扩展的场景。
  • 若需流程管控,Cloud Workflows是更好的选择;若场景简单,单函数并行处理可降低部署成本。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 03:25:26