「資料管線的終極價值,在於將散落各處的市場數據進行自動化聚合、透過 AI 進行深度解讀,並在第一時間精準推送至維運者的終端。」
在 Day 08 中,我們探討了如何透過 Google Drive API 自動同步雲端文件 (drive_sync.py)。今天我們將進入 Week 3 自動化管線的核心實踐——daily_analysis.py 每日盤後市場分析自動化腳本。
這個腳本每日定時於 RHEL 伺服器上執行,自動抓取 台灣證券交易所(TWSE) 盤後數據、Yahoo Finance 美股三大指數與財經新聞,整合本地 RAG 知識庫後,呼叫 Gemini 進行推理,並透過 Telegram Bot 將報告推送給使用者,同時記錄至 Google Sheets 試算表與存回 Angelina 知識庫中。
為了獲取全面且即時的市場全貌,daily_analysis.py 透過 httpx.AsyncClient 非同步並列處理三方數據源:
為了確保分析報告產出的高穩定度,腳本除了在單一模型提供 3 次 Exponential Backoff 重試外,還加入了 跨模型 Model Fallback 自動降級 機制:
Python
GEMINI_MODELS = ['gemini-2.5-flash', 'gemini-flash-lite-latest']
async def generate_report(market_data, knowledge):
"""使用 Gemini AI 生成市場分析報告,支援 Model Fallback 與自動重試"""
prompt = build_analysis_prompt(market_data, knowledge)
async with httpx.AsyncClient() as client:
# 第一層:模型降級輪詢 (Primary -> Fallback)
for model_name in GEMINI_MODELS:
url = f"https://generativelanguage.googleapis.com/v1beta/models/{model_name}:generateContent?key={GEMINI_API_KEY}"
payload = {
"contents": [{"parts": [{"text": prompt}]}],
"generationConfig": {"temperature": 0.7, "maxOutputTokens": 4096}
}
# 第二層:單一模型指數退避重試 (3 次)
for attempt in range(3):
resp = await client.post(url, json=payload, timeout=120)
if resp.status_code == 200:
data = resp.json()
candidates = data.get('candidates', [])
if candidates and candidates[0].get('content', {}).get('parts'):
return candidates[0]['content']['parts'][0].get('text', '')
elif resp.status_code in (503, 429):
wait = 20 * (attempt + 1)
print(f" [RETRY] {model_name} returned {resp.status_code}, waiting {wait}s")
await asyncio.sleep(wait)
else:
break # 遇到其他 4xx/5xx 直接切換至下一備援模型
print(" [ERROR] All models failed")
return None
Telegram Bot API 對單筆訊息有 4,000 字元 的嚴格限制。當生成的分析報告長度過長時,腳本會進行逐行智慧拆分,並使用 MD5 Hash 記錄於 /tmp/daily_analysis_sent_hashes.json 中,避免排程重複發送:
Python
def get_message_hash(text):
"""產生當日訊息雜湊值以進行發送去重"""
today = datetime.now(TW_TZ).strftime('%Y-%m-%d')
return hashlib.md5(f"{today}:{text.replace('\n', ' ')[:500]}".encode()).hexdigest()
async def send_telegram(msg):
"""發送報告至 Telegram,字數 >4000 時自動切分多段發送"""
msg_hash = get_message_hash(msg)
if msg_hash in load_sent_hashes():
print(" [INFO] Message already sent today, skipping duplicate.")
return True
# 訊息超過 4000 字元時按行進行拆分
messages = []
if len(msg) > 4000:
parts, current = [], ""
for line in msg.split('\n'):
if len(current) + len(line) + 1 > 4000:
parts.append(current)
current = line
else:
current = f"{current}\n{line}" if current else line
if current:
parts.append(current)
messages = parts
else:
messages = [msg]
# 依序發送 Telegram 訊息 (失敗時自動降級 Markdown -> HTML -> Plain Text)
async with httpx.AsyncClient() as client:
for i, part in enumerate(messages):
payload = {'chat_id': TELEGRAM_CHAT_ID, 'text': part, 'parse_mode': 'Markdown'}
resp = await client.post(f"https://api.telegram.org/bot{TELEGRAM_BOT_TOKEN}/sendMessage", json=payload, timeout=30)
if resp.status_code != 200:
payload['parse_mode'] = 'HTML'
resp = await client.post(url, json=payload, timeout=30)
await asyncio.sleep(1)
save_sent_hash(msg_hash)
return True
分析報告生成並推送成功後,整個工作流程會進行最後的資料閉環(Data Loop):
透過 daily_analysis.py 的實作,我們建構了完全自動化的盤後分析機器人:
明天(Day 10)我們將深入 sheets_tracker.py,探討如何透過 Service Account 與 gspread 套件建立 Google Sheets 自動化紀錄與預測準確率追蹤!
明日預告:【Day 10】數據追蹤與績效評估:Google Sheets 數據寫入與預測準確率追蹤 (sheets_tracker.py)