非同步批次處理:大量任務怎麼跑得又快又省
一萬份文件要處理,逐筆同步呼叫要跑一整天。批次架構能壓到幾小時。
批次任務與線上請求的最佳化方向相反:線上追求低延遲,批次追求高吞吐與低成本。用線上的寫法跑批次,通常慢十倍以上。
三個關鍵手段
- 併發而非序列:同時發出多個請求,但要受限流控制。
- 用批次 API:多數服務商提供離線批次介面,單價明顯較低,代價是延遲以小時計。
- 斷點續跑:一萬筆跑到第八千筆掛掉,不能從頭來。
有限併發的實作
import asyncio
async def run_all(items, worker, concurrency=8):
sem = asyncio.Semaphore(concurrency)
async def one(it):
async with sem:
return await worker(it)
return await asyncio.gather(*(one(i) for i in items),
return_exceptions=True)return_exceptions=True 很重要:一筆失敗不該讓整批中斷。失敗的項目收集起來另外重試。
斷點續跑
done = set()
if os.path.exists(OUT):
done = {json.loads(l)['id'] for l in open(OUT, encoding='utf-8')}
pending = [i for i in items if i['id'] not in done]
with open(OUT, 'a', encoding='utf-8') as f:
for result in process(pending):
f.write(json.dumps(result, ensure_ascii=False) + '\n')
f.flush() # 立刻落盤,中斷時不遺失併發數怎麼定
從服務商的每分鐘請求上限反推,並保留兩到三成餘裕給線上流量。批次任務把配額吃光導致線上服務被限流,是很常見的事故。
批次任務排在離峰時段跑。這同時降低對線上服務的影響,某些服務商的離峰配額也比較寬鬆。