iT邦幫忙

2026 iThome 鐵人賽

DAY 8
0
Build on Google AI

打造零成本企業級 AI Agent:以 Gemini 2.5 Flash 構建金融分析助手與維運實戰系列 第 8

【Day 08】知識管線自動化:Google Drive 檔案同步與 /learn 指令排程注入 (drive_sync.py)

  • 分享至 

  • xImage
  •  

「一個真正實用的 AI Agent,不該被動等待使用者手動貼入資料,而要能主動對接雲端儲存,實現知識庫的自動化增量同步。」

在先前幾天,我們完成了 AI Agent 的大腦推理、記憶與背景學習機制。從今天開始我們正式進入 資料管線自動化與周邊服務整合 的章節!

今天我們將深入探討 drive_sync.py 如何透過 Google Drive API (Service Account) 監控雲端資料夾,自動辨識新增或修改的 Google Docs、.txt 與 .md 檔案,並透過 FastAPI 的 /chat 介面搭配 /learn 特殊指令 進行動態文字切分與 Incremental Ingestion。

本篇重點摘要

  1. 為什麼選擇 Service Account 進行 Headless 雲端同步呢?
  2. 拆解 drive_sync.py 中的狀態追蹤檔(drive_sync_state.json)增量比對機制。
  3. 處理多格式檔案下載(Google Docs 導出 text/plain vs 本地 .txt/.md 直接下載)。
  4. 實作文本動態切分算法(MAX_LEARN_CHUNK = 1500)與 /learn 指令分批注入。

一、同步管線架構與安全憑證設計

在 Self-hosted RHEL 伺服器上,為了讓背景腳本能自動排程(Cron Job)同步雲端文件,我們採用了 Google Service Account(服務帳號)機制:
1. Headless Auth:無須開啟瀏覽器進行 OAuth2 使用者授權,透過 /opt/angelina/config/service-account.json 憑證檔直接取得僅讀(drive.readonly)存取 Token。
2. 輕量化 API 對接:使用 google.oauth2.service_account 搭配 httpx.AsyncClient 非同步呼叫 Google Drive v3 REST API。
3. 狀態鎖定與去重:透過記錄檔案的 modifiedTime(修改時間),確保只有新增或更動過的文件才會被下載與注入,避免重複消耗 Token。

二、drive_sync.py 核心原始碼解析

  1. 多格式文件下載與導出 (download_file_content)
    Google Drive 中的檔案類型多樣,處理時需針對不同 mimeType 採用對應的下載策略:
    (1) Google Docs (application/vnd.google-apps.document):呼叫 API 的 /export 端點導出為純文字 (text/plain)。
    (2) 純文字/Markdown (text/plain, text/markdown):呼叫 API 的 /files/{file_id} 帶入 alt=media 參數直接下載。
Python
async def download_file_content(access_token: str, file_id: str, mime_type: str) -> str | None:
    """從 Google Drive 下載或導出檔案內容"""
    headers = {"Authorization": f"Bearer {access_token}"}

    async with httpx.AsyncClient() as client:
        if mime_type == "application/vnd.google-apps.document":
            # Google Docs 導出為純文字
            url = f"{DRIVE_API_BASE}/files/{file_id}/export"
            params = {"mimeType": "text/plain"}
            response = await client.get(url, params=params, headers=headers)
            response.raise_for_status()
            return response.text

        elif mime_type in ("text/plain", "text/markdown"):
            # .txt/.md 檔案直接下載
            url = f"{DRIVE_API_BASE}/files/{file_id}"
            params = {"alt": "media"}
            response = await client.get(url, params=params, headers=headers)
            response.raise_for_status()
            return response.text

        else:
            logger.info(f"Skipping unsupported mime type: {mime_type}")
            return None
  1. 文本分塊與 /learn 指令注入 (sync_to_knowledge)
    為了避免單次 POST 請求的長度超出 API 限制,當下載的檔案字數超過 1,500 字元(MAX_LEARN_CHUNK = 1500) 時,腳本會優先在 換行符 (\n) 或 空格 處進行智慧切分,並加入分塊標籤(例如 (part 1/3)):
Python
MAX_LEARN_CHUNK = 1500

async def sync_to_knowledge(text: str, filename: str) -> None:
    """將文本切分並透過 /learn 指令寫入 Angelina 知識庫"""
    chunks = []
    # 智慧邊界切分算法 (優先搜尋 \n 或空格)
    while len(text) > MAX_LEARN_CHUNK:
        split_pos = text.rfind("\n", 0, MAX_LEARN_CHUNK)
        if split_pos == -1:
            split_pos = text.rfind(" ", 0, MAX_LEARN_CHUNK)
        if split_pos == -1:
            split_pos = MAX_LEARN_CHUNK
        chunks.append(text[:split_pos])
        text = text[split_pos:].lstrip()
    if text:
        chunks.append(text)

    async with httpx.AsyncClient(timeout=60.0) as client:
        for i, chunk in enumerate(chunks, 1):
            part_label = f" (part {i}/{len(chunks)})" if len(chunks) > 1 else ""
            # 格式化為 Angelina 內建指令格式
            message = f"/learn [Drive: {filename}{part_label}] {chunk}"
            
            response = await client.post(
                f"{ANGELINA_API_URL}/chat",
                json={"message": message, "session_id": "drive-sync", "language": "zh-TW"},
            )
            response.raise_for_status()
            logger.info(f"Ingested chunk {i}/{len(chunks)} for '{filename}'")
  1. 主排程與狀態控管 (run_sync)
    主流程會讀取本機的 /opt/angelina/data/drive_sync_state.json,對比每份檔案的 modifiedTime:
Python
async def run_sync() -> dict:
    summary = {"synced": 0, "skipped": 0, "errors": 0}
    access_token = get_drive_credentials()
    files = await list_files(access_token)
    sync_state = load_sync_state()

    for file_info in files:
        file_id, filename = file_info["id"], file_info["name"]
        modified_time = file_info["modifiedTime"]

        # 比對修改時間,未變更則跳過
        if file_id in sync_state and sync_state[file_id] == modified_time:
            summary["skipped"] += 1
            continue

        content = await download_file_content(access_token, file_id, file_info["mimeType"])
        if content:
            await sync_to_knowledge(content, filename)
            sync_state[file_id] = modified_time
            summary["synced"] += 1

    save_sync_state(sync_state)
    return summary

三、今日總結與自動化運維亮點

透過 drive_sync.py 的實作,我們完成了資料輸入端(Data Ingestion Pipeline)的自動化:

  1. 零人工介入:只要將理財筆記、研究報告放進特定的 Google Drive 資料夾,系統會在背景自動吸收到 ChromaDB 中。
  2. 複用現有架構:不需額外開闢特化 API,直接透過 /chat 端點的 /learn 指令寫入,自動觸發非同步學習與向量庫去重。
  3. 輕量高效:透過 drive_sync_state.json 比對,保證每次同步僅花費數秒,非常適合掛載於 Linux Cron Job 中運作。

明天(Day 09)我們將進入 daily_analysis.py,探討如何抓取 TWSE 證交所 API 與 Yahoo Finance 美股數據,結合 Gemini 生成每日市場分析報告並推送至 Telegram 機器人!

明日預告:【Day 09】每日自動化盤後分析:TWSE/美股數據爬取、Gemini 分析與 Telegram 機器人推播 (daily_analysis.py)


上一篇
【Day 07】非同步背景學習:LearningModule 自動提煉金融知識與向量庫去重
下一篇
【Day 09】每日自動化盤後分析:TWSE/美股數據爬取、Gemini 分析與 Telegram 機器人推播 (daily_analysis.py)
系列文
打造零成本企業級 AI Agent:以 Gemini 2.5 Flash 構建金融分析助手與維運實戰9
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言