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.")
请问该策略是否合理,是否存在更优的实现方案?
现有策略的合理性分析
你的架构思路是合理的,核心优势包括:
- 解耦性:通过Pub/Sub拆分事件ID采集与各投注市场的赔率采集逻辑,避免单一函数复杂度过高,后续新增投注市场时只需添加新的订阅函数即可,扩展性强。
- 异步执行:Pub/Sub支持异步消息传递,Function1完成采集后无需等待后续函数执行,可立即终止,提升整体流程效率。
- 故障隔离:单个投注市场的采集失败不会影响其他市场或事件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架构
关键调整点:
- 修复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()}")
- 简化订阅函数的云函数逻辑:移除本地订阅监听代码,专注处理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逻辑保持不变 # ...
- 添加错误处理与重试机制:
- 对API调用添加超时、重试(可使用
tenacity库),避免网络波动导致采集失败。 - 配置Pub/Sub订阅的死信队列,将处理失败的消息转发到死信主题,便于后续排查与重处理。
- 对API调用添加超时、重试(可使用
- 批量处理优化:若event_ids数量较多,拆分小批量进行API调用,避免单次请求参数过长被API拒绝。
方案2:使用Cloud Workflows编排流程
如果需要更严格的流程控制(例如确保所有投注市场采集完成后执行后续操作),可以用Cloud Workflows替代Pub/Sub的一对多触发,实现精细的流程编排:
- Cloud Scheduler触发Cloud Workflows。
- Workflows第一步调用Function1获取event_ids。
- Workflows并行调用Function2、Function3传递event_ids。
- 可添加流程分支、错误处理、重试逻辑,甚至在所有函数执行完成后触发通知或数据校验任务。
该方案优势是流程可视化、可追溯,适合对执行顺序和结果有明确要求的场景,但配置复杂度略高于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

