從 CSV Parsing 到 Background Data Pipeline:製造資料處理架構的演進

CASE STUDY / DATA PIPELINE

一開始我以為需求只是「把 CSV 存進資料庫」,後來才發現真正要設計的是:如何讓數量不固定、可能失敗的資料工作,在 Web Request 之外被可靠地完成,而且每一份檔案都能被追蹤。

Jupyter Notebook · pandas · Pydantic · FastAPI · Huey · Redis · diskcache · Prisma Client Python · MySQL · SMT SPI / AOI

說明:本文以製造業 MES 後端專案中的實際程式碼、notebook 與 commit history 為基礎整理。產線、站別、資料表與類別名稱皆已匿名化(例如以 LINE_A、SpiTestService 代稱),程式邏輯保持不變。文中沒有任何效能數字,因為專案裡沒有可支撐的 benchmark。


1. 前言:需求看起來只是「把 CSV 存進資料庫」

SMT 產線上的 SPI(錫膏檢測)與 AOI(自動光學檢測)設備,每檢測一片板子就輸出一份檔案。後端要做的事聽起來很單純:把這些檔案讀進來、存進 MES 資料庫,讓品保與製程工程師可以用條碼追溯每一片板子的檢測結果。

實際做下去,問題很快就不只是 parsing:

  • 每家設備的 CSV 都不是「一張表」,而是好幾個區塊拼在一起的報表格式
  • 一份檔案展開後是 pad-level 的明細,一次寫入就是一批資料
  • 資料要跨越「產線電腦 → API → 資料庫」三個邊界,任何一段都可能失敗
  • 新站別持續加入,parsing rule 會一直改

這篇記錄的是這條 pipeline 如何從 notebook 實驗、Client 端解析,演進成「上傳檔案 → Server 解析驗證 → Queue → Worker 寫入 → 狀態查詢」的架構,以及每一步背後的取捨。

2. 從 Notebook 開始:先搞懂資料長什麼樣

我在 repository 裡替每一種設備格式都留了一本日期命名的 notebook(例如 SPI、AOI、OS、效率測試、HIPOT 各一本)。它們的角色不是教學,而是在寫任何 API 之前,先把未知格式的 parsing rule 定下來。

2.1 一份 CSV 其實是三個區塊

以某台 SPI 設備為例,notebook 第一步只是把原始內容印出來看:

生产线名,Line01
机器名称,<machine>
操作员姓名,<operator>
检查程序名,<program>-B
批号,<lot>
MOUNT块编号,0,1,2,3,
条形码,<barcode-1>,<barcode-2>,<barcode-3>,<barcode-4>,
检查焊盘数,1376
Ng数,0
WARN数,1
判定结果,程序块编号,MOUNT块编号,...,NG上限,NG下限,...,NG上限,NG下限,...
OK,0,0,...,  134.6,210,40,...

從這裡可以直接歸納出幾條規則:

  • 檔頭是 key-value,但被中間兩行條碼切成兩段(第 1–5 行、第 8–10 行)
  • 條碼是橫向排列,一列 mount block 編號、一列條碼,兩者要一一對應
  • 明細表從第 11 行開始,而且欄位名稱大量重複(NG上限 出現九次,分別屬於體積、高度、面積……)
  • 數值帶前導空白(' 104.8'),- 代表沒有值
  • 編碼因設備而異:這台是 GBK、簡體中文欄位;另一條線的 SPI 是 UTF-8、英文欄位
  • 部分 metadata 只存在檔名裡:<code>_<YYYYMMDDhhmmss>_<OK|NG|WN|Judge>.csv

2.2 把規則寫成可重用的函式

notebook 裡的 parsing function 是按區塊拆開的,這個結構後來幾乎原封不動搬進正式系統:

def read_spi_header(file_path: str) -> dict:
    # 檔頭是 key-value,但被條碼區塊切成兩段
    df_1_5 = pd.read_csv(file_path, header=None, skiprows=0, nrows=5,
                         usecols=[0, 1], encoding="gbk")
    df_8_10 = pd.read_csv(file_path, header=None, skiprows=7, nrows=3,
                          usecols=[0, 1], encoding="gbk")

    df = pd.concat([df_1_5, df_8_10], ignore_index=True)
    df.fillna("", inplace=True)
    df.columns = ["headers", "data"]
    return df.set_index("headers")["data"].to_dict()


def parse_filename(filename: str):
    pattern = re.compile(r"^(?P<code>.+)_(?P<ts>\d{14})_(?P<judge>OK|NG|WN|Judge)$")
    m = pattern.match(filename)
    return (m["ts"], m["judge"]) if m else None

重複欄位名的問題,pandas 讀取時會自動改成 NG上限、NG上限.1、NG上限.2……所以我用一張 mapping table 把「位置語意」轉成「欄位語意」,例如 NG上限.1 → height_ng_upper。型別轉換則交給 Pydantic:str_strip_whitespace=True 處理前導空白,model_validator(mode="before") 統一把字串轉成 int / float。

2.3 Notebook 裡的實驗演進

幾個從 notebook 看得出來的迭代:

  • 讀一次再切片:早期同一份檔案用 skiprows/nrows 讀兩三次;後來另一條線的版本改成 read_csv 讀一次,再用 iloc[0:4]、iloc[6:9] 切出需要的列。正式程式碼採用後者。
  • 編碼偵測:在另一本 FCT 的 notebook 裡試過 chardet 偵測編碼,正式程式碼則是依設備寫死編碼(GBK / UTF-8)。
  • 一個被吞掉的錯誤:讀條碼的函式包了 try/except 回傳 (None, None)。notebook 輸出顯示某些檔案確實拿到了 (None, None)——這代表 fallback 會把「格式錯誤」和「這片板子沒有條碼」混成同一種結果。這個問題到現在的版本仍然存在,後面會再提。

Notebook 的價值在於: 在碰 API、交易、佇列之前,先用最短的回饋迴圈回答「這份資料的結構是什麼、哪些欄位可信、哪些值要正規化」。

3. 第一版:在 Client 端解析資料

parsing rule 穩定之後,第一個上線的做法是在產線電腦上跑 Python script(repository 裡的 *_csv2db.py):

  1. 從 config.ini 讀取要監看的資料夾
  2. 逐一掃描 CSV,用 notebook 裡同一套 pandas 函式解析
  3. 在 Client 端建立 Pydantic model(Server 端也有一份同樣的 model)
  4. 把結果以 JSON POST 到 API
  5. 上傳成功就把檔案搬到 uploaded/ 資料夾

較早的站別(其他產線的 FCT、ICT、鎖螺絲、點膠、效率測試等)都是這個模式。SPI 則拆成了兩個 API:先建立檢測主檔拿回 test_id,再把明細分批送出。(拆分的原因根據目前 repository 推測是明細筆數多——notebook 範例檔的檔頭記錄的檢查焊盤數在一千到兩千多之間,明細是 pad-level 的列。)

spi_response = await upload_data_async(spi_data.model_dump())  # POST /SPI_Test
test_id = spi_response.json().get("id")
if test_id is None:
    logger.info(f"Can not upload file {file_name}")
    continue

await upload_data_items_async(test_id, item_list)   # POST /SPI_Test/{test_id},每批 20 筆
move_uploaded_file(file_path, uploaded_dir)
Before:產線電腦上的 csv2db script 以 pandas 解析 CSV、建立 Client 端 Pydantic model,先 POST 主檔取得 test_id,再分批 POST 明細到 /SPI_Test/{test_id},最後寫入 MES DB;全部成功後才把檔案搬到 uploaded/
Before:解析在產線電腦上完成,一份 CSV 被拆成「主檔+多批明細」多個 request 寫入。

這個版本一開始的優點很實際:

  • Server 端只要提供 CRUD API,實作快
  • notebook 的程式碼可以直接複製到 script,驗證過的 parsing rule 立刻能用
  • 每個新站別都能獨立交付,不影響既有 API

4. 問題:Parsing 不應該綁在 Client

用了一陣子之後,Client 端的 log 暴露出這個架構的邊界問題。以下都是 log 或 commit 裡看得到的事實,不是事後推論。

4.1 一份檔案被拆成多個 request,寫入不是原子的

主檔 API 成功、但後續某幾批明細失敗(log 裡有 422、500、連線失敗)時,資料庫裡就會留下一筆「有主檔、明細不完整」的檢測紀錄。Server 端為此補了一段補償邏輯:寫入失敗時先刪掉該檔案的條碼與明細再重送。

這是一個訊號:「一份 CSV = 一個業務單位」這個不變式,被 HTTP request 的切法打破了,而維護不變式的責任落在一支跑在產線電腦上的 script。

4.2 Schema 有兩份,而且會漂移

Client 和 Server 各有一份同名的 Pydantic model。以目前 repository 的狀態來看,兩份已經不一致:Client script 裡的 model 把條碼定義成逗號分隔的 str、檢測時間也是 str;Server 端的同名 model 則是 list 與 datetime。Client 的 log 裡也出現過整批 Input should be a valid string 的 validation error——parsing 輸出的型別調整了,model 卻沒跟上。

資料格式、parsing rule、schema 的所有權分散在兩個部署單位,任何一邊改動都要同步另一邊。

4.3 錯誤處理、重試與可觀測性都在 Client

失敗時要不要重試、檔案要不要搬、錯誤記在哪裡,全部寫在 script 裡;log 也只存在那台電腦上。Server 看不到「有一份檔案處理失敗」,只看得到零散的 request error。

4.4 所以要移動的是 execution boundary

問題不在於「背景處理比較好」,而在於誰擁有這份資料的語意。parsing rule、schema、驗證與寫入的一致性,都應該在同一個地方被定義和執行。Client 該做的事情只剩一件:可靠地把原始檔案送到 Server。

Client 只負責把原始檔案可靠地送到 Server;資料的語意只由 Server 擁有。

5. 第二版:上傳檔案,由 Server 負責其餘一切

新架構下,Client script 被簡化成純粹的 uploader:

async def upload_file(client, filename_with_ext: str, file_bytes: bytes):
    return await client.post(
        url=f"{BASE_URL}/LINE_A/SPI_Test/file-upload",
        files={"file": (filename_with_ext, file_bytes, "text/csv")},
    )

# 上傳前先問 Server 哪些檔案已經存在,只送新檔
existing_files = await get_existing_filenames(client)   # GET /SPI_Test/file-names
new_files = [f for f in files if stem(f) not in existing_files and stem(f) not in error_files]

它還保留了幾個「好公民」的設定:每批 5 個檔案、同時只上傳 1 個、每批之間休息 1 秒,失敗的檔案複製到 error/ 資料夾,下次掃描時略過。

這個轉變不是一步到位的。從 commit history 看得到三個階段:

階段解析在哪裡寫入 DB 在哪裡API 回應
A. Client 解析產線電腦request 內寫入結果
B. 檔案上傳(同步)Server,request 內request 內寫入結果
C. 檔案上傳 + 背景寫入Server,request 內Huey workerqueued + task_id

部分站別(例如某台輸出 .txt 的 SPI、以及電性測試站)目前仍停在階段 B;兩條產線的 SPI 與 AOI 共四個 endpoint 已經走到階段 C。

目前的完整架構如下:

After:檢測設備 CSV 由 Uploader script 以 multipart 上傳到 FastAPI file-upload endpoint,在 request 內 Decode、Parse、Validate(格式不符回 422),再以 enqueue_upload_task 放入 Redis Huey queue 並在 diskcache 記錄 UploadCsvTask 狀態;Huey consumer 透過 Service 與 Unit of Work 寫入 MES DB(MySQL,Prisma Client),signals 更新狀態,Status API 查詢狀態
After:Client 只上傳原始檔案;解析驗證在 Server 的 request 內完成,資料庫寫入交給 Huey worker,狀態可由 Status API 查詢。

資料的生命週期:

  1. Uploader 送出原始檔案(bytes)
  2. API 依設備編碼 decode,解析檔名與三個區塊,組成 Pydantic model——檔案不合法就在這裡直接回 422
  3. 驗證過的 model 被放進 Redis queue,同時在 diskcache 建立一筆 PENDING 的 task 狀態
  4. API 回傳 task_id,request 結束
  5. Worker 取出 task,開 transaction 寫入主檔、條碼、明細
  6. Huey signal 把狀態更新成 STARTED → SUCCESS / FAILURE(重試中會回到 PENDING)
  7. 使用者或 script 用 task_id 查詢處理結果

6. Upload Endpoint:接收成功 ≠ 處理完成

@router.post("/SPI_Test/file-upload", status_code=status.HTTP_201_CREATED)
async def upload_spi_test_file(
    file: UploadFile = File(...),
    spi_service: SpiTestService = Depends(create_spi_test_service),
):
    file_name = spi_service.split_file_name_and_check_ext(file.filename)  # 非 .csv → 422

    # 同步:decode + parse + validate,錯誤在這裡就回給 client
    spi_test_model, item_models = spi_service.process_upload_and_validate(file_name, file.file)
    spi_test_model.items = item_models

    # 非同步:寫入 DB 交給 worker
    task = enqueue_upload_task(upload_spi_data, file_name, spi_test_model, "SPI")

    return {"status": "queued", "filename": file_name, "task_id": task.id}

這個 endpoint 的責任邊界很明確:

  • 同步做的事:副檔名檢查、解碼、parsing、schema 驗證。這些是「這份檔案本身有沒有問題」,應該立刻讓上傳端知道,而不是排進佇列才失敗。
  • 非同步做的事:資料庫寫入。這是「系統有沒有辦法把它存好」,可能受資料庫狀態、連線、鎖影響,適合交給有重試機制的 worker。

API 接收成功,不等於資料處理完成。

回應裡的 "status": "queued" 與 task_id,就是在 API contract 上明確表達「我收到了,而且格式正確,但還沒寫完」。

老實說,這裡的 HTTP status code 還是 201 Created,語意上 202 Accepted 更準確——同一個專案裡另一個背景匯出任務的 endpoint 就是用 202。這是之後應該統一的地方。

7. Background Worker:Huey + Redis

選 Huey 不是憑空決定的。專案裡早就有一個用 Huey(SQLite backend)做的背景任務:把大量條碼的檢測資料打包成 zip 供下載,當時也留有一本 Huey 的實驗 notebook。CSV 上傳沿用同一套工具,但換成 RedisHuey,用獨立的 queue name 和舊任務隔開。

7.1 Enqueue 與狀態追蹤

upload_huey = RedisHuey("upload_tasks", host=REDIS_HOST, port=6379)
upload_task_cache = Cache(directory="./cache/upload-task-queue/")

def enqueue_upload_task(task_wrapper, file_name: str, data_model, task_name: str) -> UploadCsvTask:
    result = task_wrapper(file_name, data_model)          # 放進 Redis,拿回 task id
    task = UploadCsvTask(
        id=result.id,
        task_name=task_name,
        file_name=file_name,
        detail_count=len(getattr(data_model, "items", []) or []),
        barcode_count=_count_barcodes(getattr(data_model, "barcodes", None)),
    )
    upload_task_cache.set(result.id, task, expire=_get_upload_task_expire())
    return task

Huey 的 result 被讀取後就會從 queue 移除,因此 task 狀態另外存在 diskcache,並設定保留期限(開發環境 10 分鐘、其他環境 10 天),讓使用者事後仍能查詢。

7.2 Worker 端:同步 worker 呼叫 async service

Service 與 repository 都是 async(Prisma Client Python 的 asyncio interface),但 Huey consumer 執行的是一般同步函式,所以 task 內部要自己驅動 event loop:

loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)

@upload_huey.task(retries=3, retry_delay=10)
def upload_spi_data(file_name: str, spi_test_model):
    spi_service = create_spi_test_service()

    async def run_task():
        await spi_service.add_spi_test_data(file_name, spi_test_model)

    return loop.run_until_complete(run_task())

這段有一個真實的演進故事。第一版的 task 長這樣:

@upload_huey.task()
def upload_spi_data(file_name, spi_test_model):
    spi_service = create_spi_test_service()
    spi_service.add_spi_test_data(file_name, spi_test_model)   # async 函式,沒有 await

呼叫 async 函式卻沒有 await,只會產生一個 coroutine 物件,寫入根本不會發生。同一版的 endpoint 還同時保留了同步寫入,所以功能看起來是正常的——資料其實是 request 裡那次寫進去的。幾天後我把同步寫入註解掉、改成上面的 event loop wrapper,並加上 retries=3, retry_delay=10。這是一個很典型的教訓:背景任務必須能被單獨驗證,不能讓同步路徑掩蓋它有沒有真的執行。

7.3 用 Signal 維護狀態機

@upload_huey.signal()
def handle_upload_task_signal(signal: str, task, exc: Exception | None = None):
    upload_task = upload_task_cache.get(task.id)
    if not upload_task:
        logger.error({"message": "Upload task ID not found in cache DB", "task_id": task.id})
        return

    if signal == signals.SIGNAL_EXECUTING:
        upload_task.state, upload_task.started_at = "STARTED", datetime.now()
    elif signal == signals.SIGNAL_COMPLETE:
        upload_task.state, upload_task.finished_at = "SUCCESS", datetime.now()
        upload_task.exception = None
    elif signal == signals.SIGNAL_ERROR:
        upload_task.state, upload_task.finished_at = "FAILURE", datetime.now()
        upload_task.exception = repr(exc) if exc else None
    elif signal == getattr(signals, "SIGNAL_RETRYING", None):
        upload_task.state = "PENDING"
        upload_task.exception = repr(exc) if exc else upload_task.exception

    upload_task_cache.set(task.id, upload_task, expire=_get_upload_task_expire())

對外提供兩個查詢端點:GET /task-queue/upload-csv/{task_id} 查單筆,GET /task-queue/upload-csv/list 可依 state、task_name、file_name 過濾。狀態裡還帶了 detail_count 和 barcode_count,出問題時可以先判斷是不是「大檔」。

8. CSV Parsing Pipeline

Server 端的 parsing 分成三層:FileReader(infrastructure) 負責讀格式、Service 負責 mapping 與組裝、Pydantic model 負責型別與驗證。

Server 端 parsing pipeline:Raw bytes 依設備以 GBK 或 UTF-8 decode 成 StringIO buffer,拆成檔頭 key-value、條碼區塊、明細表;檔頭與明細經 key mapping 轉成語意欄位,與條碼區塊、檔名 regex 解析結果一起經 Pydantic 驗證,組成主檔加 items 的 Create DTO
Server 端 parsing:Decode → 區塊切分 → Key mapping → Pydantic 驗證 → Create DTO。
def process_upload(self, filename: str, file: BinaryIO):
    test_date_time, check_result = self.file_reader.parse_filename(filename)
    content = file.read().decode("gbk")          # 這台設備固定輸出 GBK
    csv_buffer = StringIO(content)

    spi_test = self.file_reader.read_spi_test(csv_buffer)
    csv_buffer.seek(0)
    spi_test_items = self.file_reader.read_spi_test_item(csv_buffer)
    csv_buffer.seek(0)
    barcodes = self.file_reader.read_spi_barcode(csv_buffer)

    return spi_test, spi_test_items, barcodes, test_date_time, check_result

def map_dict_to_model(self, raw: dict, mapping: dict) -> dict:
    result = {}
    for k, v in raw.items():
        if k in mapping:
            result[mapping[k]] = None if isinstance(v, str) and v.strip() == "" else v
    return result

幾個設計點:

  • 同一個 buffer 讀三次:上傳進來的是 SpooledTemporaryFile,先整份 decode 成字串放進 StringIO,每個區塊讀完 seek(0) 再讀下一個,避免依賴底層檔案物件的狀態。
  • mapping table 是資料,不是程式邏輯:設備改欄位名時只需要改 dict。
  • 不同設備、不同 reader:SPI 的「檔頭+條碼+明細」、AOI 的「寬表(60+ 欄)+固定欄位位置的條碼」各自有 FileReader,但都輸出相同形狀的 (主檔 dict, 明細 list),讓 Service 與 Repository 不需要知道格式差異。
  • 板面(Top / Bottom)不在資料裡:SPI 從程式名稱結尾推得,AOI 從檔名的料號段推得,這些規則都是在 notebook 裡先驗證過的。

9. Database Persistence

資料庫存取使用 Prisma Client Python(asyncio)連到 MES 的 MySQL,外面包一層 Unit of Work,把 Prisma 的 tx() 包成 async context manager:

class SpiTestUnitOfWork:
    def __init__(self):
        self.tx_manager = mes_db.tx()
        self.spi_test_repo: SpiTestRepository | None = None

    async def __aenter__(self):
        if not mes_db.is_connected():
            await mes_db.connect()
        tx = await self.tx_manager.__aenter__()
        self.spi_test_repo = SpiTestRepository(tx)
        return self

    async def __aexit__(self, exc_type, exc_val, exc_tb):
        await self.tx_manager.__aexit__(exc_type, exc_val, exc_tb)

寫入流程與重複處理:

async def add_spi_test_data(self, filename: str, data: SpiTestCreate):
    try:
        spi_data, barcodes, mount_blocks, items = await self._split_data(data)

        async with SpiTestUnitOfWork() as uow:                 # tx 1:主檔 + 條碼
            spi_test = await uow.spi_test_repo.add_one_test(spi_data)
            await uow.spi_test_repo.add_barcodes(spi_test.id, barcodes, mount_blocks)

        for batch in self._batch(items, self.BATCH_SIZE):       # tx 2..N:明細,每批 10 筆
            async with SpiTestUnitOfWork() as uow:
                await uow.spi_test_repo.add_items(spi_test.id, [i.model_dump() for i in batch])
        return spi_test

    except EntityAlreadyExistError:                             # 同一份檔案再次進來
        spi_test = await self.query_repo.find_test_by_file_name(filename)
        async with SpiTestUnitOfWork() as uow:
            await uow.spi_test_repo.remove_barcodes(spi_test.id)
            await uow.spi_test_repo.remove_items(spi_test.id)
            await uow.spi_test_repo.add_barcodes(spi_test.id, barcodes, mount_blocks)
        for batch in self._batch(items, self.BATCH_SIZE):
            async with SpiTestUnitOfWork() as uow:
                await uow.spi_test_repo.add_items(spi_test.id, [i.model_dump() for i in batch])
        return spi_test

Repository 把資料庫例外翻譯成 domain exception:

async def add_one_test(self, data: dict):
    try:
        return await self.spi_test.create(data)
    except UniqueViolationError:
        raise EntityAlreadyExistError(message=f"The test of file [{data['file_name']}] is already existed")

這裡最重要的設計是 file_name 上的 unique constraint。它讓「同一份檔案」成為資料庫層級可辨識的 identity,因此:

  • Uploader 可以用 file-names API 做上傳前去重
  • 重複上傳或 worker 重試時,主檔不會重複建立,而是走「刪除舊明細 → 重新寫入」的替換路徑
  • 搭配 Huey 的 retries=3,一個寫到一半失敗的 task 在重試時會收斂到完整狀態,而不是留下兩份資料

也要誠實說明目前的取捨:

  • 不是單一 transaction:主檔與每一批明細是分開的 transaction。好處是單一 transaction 不會過大;代價是失敗時資料庫會短暫存在「不完整」的中間狀態,依賴重試來修復。
  • 明細是逐筆寫入:repository 裡保留了被註解掉的 create_many 版本,目前實際執行的是在 transaction 內迴圈逐筆寫入。
  • 替換路徑只換明細:重送同名檔案時,條碼與明細會被替換,但主檔欄位不會更新。如果設備重新輸出同名但內容不同的檔案,主檔會停留在舊值。

10. Failure Handling:Data Pipeline 不只是 Happy Path

這個專案目前沒有針對 CSV pipeline 的自動化測試(tests/ 裡的測試涵蓋的是其他模組)。能證明這條 pipeline 如何處理失敗的,是 notebook 的實驗輸出、Client 端 log,以及程式碼本身。以下依失敗發生的位置整理:

失敗情境目前行為依據
副檔名不是 .csv同步回 422(UnSupportedExtError),不進 queueService 程式碼
欄位缺漏、型別錯誤Pydantic 驗證失敗,request 內就失敗Model 程式碼、Client log 中的 422
檔名不符合 regexparse_filename 回傳 None,解包時拋出非 domain 的例外程式碼推論
條碼區塊格式異常被 except 吞掉,回傳 None,資料照樣寫入但沒有條碼FileReader 程式碼、notebook 輸出
寫入明細途中 DB 失敗Huey 重試最多 3 次,每次間隔 10 秒;重試時走替換路徑retries=3、EntityAlreadyExistError 分支
重試用盡task 狀態變成 FAILURE,exception 欄位保留 repr(exc)signal handler
同一份檔案重複上傳上傳端先去重;即使送到,也只替換明細file_name @unique、file-names API

從這張表可以看出一個刻意的切分:格式問題在 request 內 fail fast,基礎設施問題交給 worker 重試。

同時也有幾個我認為目前版本仍存在的限制:

  • 原始檔案沒有保留:Server 讀完 UploadFile 就丟掉,queue 裡放的是解析後的 Pydantic model(由 Huey 序列化存入 Redis)。這表示無法在修正 parser 後重新處理舊檔,除錯時也看不到原始輸入。
  • Worker 當機時的行為取決於 Huey:還在 Redis 裡排隊的 task 會保留;但根據 Huey 的設計,task 在執行前就會從 queue 取出,執行到一半時 consumer 被砍掉的那一筆,不保證會被重新執行。
  • 狀態存在本機的 diskcache:API 與 worker 必須共用同一個檔案系統才看得到同一份狀態,這限制了兩者的部署方式。另外 enqueue_upload_task 是先 enqueue 再寫入狀態,根據程式碼推測,worker 動作夠快時可能先收到 EXECUTING signal 而找不到狀態紀錄(handler 會記一筆 error log 後略過)。
  • 靜默的型別轉換:數值轉換失敗時預設回傳 0 / 0.0,而不是拒絕該筆資料。對檢測資料而言,「0」和「缺值」意義不同,這是一個資料完整性風險。

11. Before vs After

Local Parsing(Client 解析)Server-side Pipeline(目前版本)
Client 責任解析、驗證、組 JSON、分批上傳、維護 schema 副本、處理部分失敗掃描資料夾、去重、節流、上傳原始檔案
Payload解析後的 JSON;SPI 需拆成主檔+多批明細 request單一 multipart 檔案
一份檔案的寫入跨多個 HTTP request,任一失敗就留下半筆資料Server 內部處理;主檔與明細分批 transaction,由重試收斂
Retry由 script 自行判斷Huey retries=3, retry_delay=10,搭配 unique file_name 的替換路徑
Failure isolation錯誤散落在各台產線電腦的 log格式錯誤同步回 422;寫入錯誤留在 task 狀態與 Server log
Parsing consistencyClient 與 Server 各有一份 model,曾出現 schema 漂移parsing rule 與 schema 只存在 Server
可觀測性只有 Client 本機 logtask_id 查詢、狀態、detail_count、例外字串
Request lifecycle等待所有明細寫完才結束解析驗證完即回應,不必等待資料庫寫入
擴展方式每個新站別要寫新 script 並部署到產線電腦新增 FileReader + Service + task,uploader 幾乎可重用

12. 如果重新設計一次(Future Improvement)

以下都是尚未實作的改進方向,依我認為的優先順序排列:

  1. 先保存原始檔,再處理(Future Improvement):上傳後先寫入 staging 目錄或 object storage,job 只帶檔案路徑與 hash。這能換來 reprocessing、audit、除錯與較小的 queue payload,也讓 parsing 可以選擇性移到 worker。
  2. 以內容 hash 作為 idempotency key(Future Improvement):目前 identity 是檔名。加入 content hash 後,可以區分「同一份檔案重送」與「同名但內容不同」,並正確更新主檔。
  3. 一份檔案一個 transaction,或 staging table 切換(Future Improvement):用 create_many 在單一 transaction 寫入,或先寫 staging 再原子切換,讓資料庫不再出現中間狀態。
  4. 把 task 狀態搬進資料庫(Future Improvement):以 PENDING / STARTED / SUCCESS / FAILURE / DEAD 的 state machine 存在 DB,解除 API 與 worker 必須共用檔案系統的限制,並消除 enqueue 與狀態寫入的競態。
  5. Dead Letter Queue 與通知(Future Improvement):重試用盡的 task 進入 DLQ,搭配專案既有的通知機制告警。
  6. OpenTelemetry 延伸到 worker(Future Improvement):FastAPI 端已經有 OTel instrumentation,下一步是讓 trace context 跟著 task 走,串起「上傳 → 排隊 → 寫入」的完整 span。
  7. Parser contract tests(Future Improvement):把 notebook 裡用過的樣本檔(去識別化後)變成 fixture,對每台設備的 FileReader 寫 golden test,涵蓋編碼、重複欄位、空值、缺條碼、錯誤檔名。
  8. API 語意與 schema versioning(Future Improvement):佇列型 endpoint 統一回 202 Accepted 並附上狀態查詢位置;在 model 中記錄設備格式版本,讓格式變更可追溯。

13. 結論

一開始,我把這個需求理解成 CSV parsing:搞懂格式、寫好 pandas、把資料塞進資料庫。

做到後來才發現,真正需要設計的是:如何讓一批數量不固定、可能失敗、處理時間不固定的資料工作,在 Web Request 之外被可靠地完成,並且讓每一份檔案的結果都可以被追蹤。

這條 pipeline 的演進,本質上是在移動責任邊界:

  • 把 parsing rule 與 schema 從產線電腦收回 Server,讓資料語意只有一個擁有者
  • 把「格式是否正確」留在 request 內 fail fast,把「能否寫入成功」交給有重試的 worker
  • 用資料庫的 unique constraint 定義檔案的 identity,讓重送與重試可以收斂
  • 用 task_id 與狀態機,讓「收到了」和「做完了」在 API contract 上成為兩件不同的事

它還不是終點——原始檔保留、單一 transaction、持久化的狀態機都還沒做。但從 notebook 到現在的 Upload → Parse & Validate → Queue → Persist → Status,這是這個專案裡我認為最重要的一次 architecture evolution。


Tech Stack

Python · FastAPI · pandas · Pydantic · Huey · Redis · diskcache · Prisma Client Python · MySQL · httpx · Jupyter Notebook · loguru