據(jù)管道最后一百米的格式轉(zhuǎn)換實(shí)戰(zhàn))
簡(jiǎn)介面向Project Diablo 2PD2玩家與腳本開(kāi)發(fā)者的Kolbot機(jī)器人腳本合集主要解決游戲自動(dòng)化操作與代理管理問(wèn)題適合已有D2BS基礎(chǔ)、希望自定義機(jī)器人行為的初、中級(jí)用戶。壓縮包共229個(gè)文件約707KB主要包含151個(gè)JavaScript腳本核心業(yè)務(wù)邏輯、33個(gè)txt配置文檔、22個(gè)nip物品拾取過(guò)濾規(guī)則、10個(gè)dbj任務(wù)啟動(dòng)文件目錄劃分明確便于按需定位和修改。目前已有256人學(xué)習(xí)/瀏覽。腳本中內(nèi)置多項(xiàng)實(shí)用配置與排錯(cuò)指引可在OOG.js第6行修改gameserver參數(shù)以指定GS服務(wù)器提供技能ID查詢指引、Kolbot NIP文件抓取配置指南還整理了D2BS崩潰的常見(jiàn)修復(fù)方法例如更新PD2BS、為D2Bot.exe和game.exe設(shè)置管理員權(quán)限能有效降低腳本部署與運(yùn)行時(shí)的排錯(cuò)成本尤其是需要頻繁調(diào)整拾取策略或服務(wù)器設(shè)定的場(chǎng)景實(shí)用性更強(qiáng)。整個(gè)包體雖小但注釋與文檔較完整適合邊用邊學(xué)。1. pd2bs-scripts 到底解決什么問(wèn)題數(shù)據(jù)管道最后一百米的格式轉(zhuǎn)換早上七點(diǎn)的定時(shí)任務(wù)打印了一屏紅色堆棧下游業(yè)務(wù)系統(tǒng) BS 拒絕了一整批訂單數(shù)據(jù)。排到中午才發(fā)現(xiàn)不是網(wǎng)絡(luò)問(wèn)題而是上游導(dǎo)出的金額還是元BS 接口只要分時(shí)間還是本地格式BS 接口要求帶時(shí)區(qū)的 UTC 字符串。這種「數(shù)據(jù)管道最后一百米」的格式適配就是 pd2bs-scripts 這類腳本存在的理由。pd2bs 是 Pipeline Data to Business System 的縮寫(xiě)pd2bs-scripts 是一套把上游管道產(chǎn)出數(shù)據(jù)PD轉(zhuǎn)換成下游業(yè)務(wù)系統(tǒng)BS可消費(fèi)報(bào)文的腳本集合。它不負(fù)責(zé)傳輸和存儲(chǔ)只負(fù)責(zé)把數(shù)據(jù)變成下游接口認(rèn)識(shí)的樣子。適合讀這篇的人是每天和批量導(dǎo)入、系統(tǒng)間數(shù)據(jù)搬運(yùn)打交道的后端或數(shù)據(jù)開(kāi)發(fā)。接下來(lái)我從數(shù)據(jù)形態(tài)、最小腳本、參數(shù)調(diào)優(yōu)講到真實(shí)翻車(chē)記錄把整條鏈路完整拆開(kāi)。2. 拆解 PD 與 BS轉(zhuǎn)換鏈路里必須先看明白的兩個(gè)邊界2.1 PD 數(shù)據(jù)長(zhǎng)什么樣JSON Lines 與字段漂移上游管道每天凌晨導(dǎo)出訂單落到共享目錄或?qū)ο蟠鎯?chǔ)文件名帶日期內(nèi)容是一行一筆訂單的 JSON Lines。我見(jiàn)過(guò)的最典型樣例長(zhǎng)這樣{order_id:20240518-10293,user_id:8899123,sku:SKU-A1,num:2,amount:98.50,paid_at:2024-05-18 03:22:11,status:PAID,extra:{coupon:C-12}} {order_id:20240518-10294,user_id:8899124,sku:SKU-B7,num:1,amount:198.00,paid_at:2024-05-18 03:25:47,status:PENDING,extra:{}}選 JSON Lines 而不是一整份大 JSON是因?yàn)樗梢宰芳?、可以按行斷點(diǎn)續(xù)讀某一行解析失敗不影響其他行。但代價(jià)是字段約束基本靠自覺(jué)上游加一個(gè)字段、改一個(gè)枚舉值下游完全不知道。這就是字段漂移。我第一次對(duì)接時(shí)按文檔寫(xiě)好了解析結(jié)果上線當(dāng)天就遇到一行業(yè)務(wù)新加的refund_time雖然不影響解析但提醒我一個(gè)事實(shí)——PD 的格式不是不能變而是變了之后必須有人負(fù)責(zé)兜住。做 pd2bs 前的第一件事不是寫(xiě)代碼而是把上游導(dǎo)出目錄里最近三天的文件都拉下來(lái)逐行數(shù)一遍字段記錄哪些字段出現(xiàn)過(guò)、哪些字段有時(shí)缺失、哪些字段的值域比文檔寫(xiě)的更寬。這個(gè)動(dòng)作花不了二十分鐘但能省掉后面大部分瞎猜。2.2 BS 接口的約束契約比想象中嚴(yán)格下游 BS 系統(tǒng)的批量接口文檔通常不長(zhǎng)但每個(gè)字段都有講究。我這邊要對(duì)接的接口長(zhǎng)這樣{ service_code: order_sync, batch_id: pd2bs_20240518_001, items: [ { outer_id: 20240518-10293, user_id: 8899123, sku_code: SKU-A1, quantity: 2, amount_cents: 9850, paid_time: 2024-05-18T03:22:11Z, status_code: 1 } ] }注意幾個(gè)和 PD 數(shù)據(jù)的差異金額從元變成分而且是整數(shù)時(shí)間從無(wú)時(shí)區(qū)的本地時(shí)間變成帶 Z 的 UTC ISO8601狀態(tài)從字符串枚舉變成數(shù)字枚舉order_id改名outer_id。每一處差異都是一個(gè)小坑合起來(lái)就是「為什么不能直接把上游文件轉(zhuǎn)發(fā)給下游」的答案。拿到接口后的標(biāo)準(zhǔn)動(dòng)作是把字段約束抄成一張對(duì)照表然后逐字段核對(duì)上游樣例數(shù)據(jù)PD 字段BS 字段類型差異轉(zhuǎn)換規(guī)則是否必填order_idouter_id字符串→字符串原樣透?jìng)魇莡ser_iduser_id字符串→整數(shù)去前導(dǎo)零后轉(zhuǎn) int是skusku_code字符串→字符串原樣透?jìng)魇莕umquantity整數(shù)→整數(shù)原樣透?jìng)餍?0是amountamount_cents字符串→整數(shù)元轉(zhuǎn)分杜絕浮點(diǎn)是paid_atpaid_time字符串→字符串本地時(shí)區(qū)轉(zhuǎn) UTC ISO8601是statusstatus_code字符串→整數(shù)PAID→1, REFUNDED→2, PENDING→3是這張表就是后面映射配置的原型。我一般會(huì)把它直接寫(xiě)成注釋掛在映射配置文件頂部因?yàn)榘肽旰蠡貋?lái)看腳本的人往往就是我自己而我最需要的恰恰是當(dāng)初核對(duì)過(guò)什么、為什么這樣映射。2.3 為什么中間必須有一層腳本直接在管道里改的三個(gè)問(wèn)題有人會(huì)問(wèn)既然差異這么明確讓上游管道在導(dǎo)出時(shí)就按 BS 的格式輸出不就行了理論上可以實(shí)操中幾乎走不通。我見(jiàn)過(guò)太多團(tuán)隊(duì)試圖這么干最后都退了回來(lái)原因有三個(gè)。第一上游管道不是只有 BS 一個(gè)下游。它要給對(duì)賬系統(tǒng)、數(shù)倉(cāng)、報(bào)表各導(dǎo)一份格式是多方博弈后的平衡。為了一個(gè)下游的需求改動(dòng)通用導(dǎo)出邏輯需要所有下游一起回歸測(cè)試周期以周計(jì)。第二映射規(guī)則變化太快。BS 接口升級(jí)、狀態(tài)枚舉調(diào)整、新業(yè)務(wù)字段接入這些都是按月出現(xiàn)的需求。如果映射邏輯燒在管道代碼里每次調(diào)整都要走發(fā)布流程。第三管道任務(wù)沒(méi)有兜錯(cuò)位置。轉(zhuǎn)換失敗的數(shù)據(jù)需要停下來(lái)給人看而不是混在管道日志里被滾動(dòng)沖掉。所以常見(jiàn)做法是讓上游只負(fù)責(zé)「把數(shù)據(jù)導(dǎo)出來(lái)」所有格式適配下沉到腳本層。pd2bs 就是這一層的實(shí)現(xiàn)輸入是上游文件輸出是 BS 接口報(bào)文中間的一切變化都在可控范圍內(nèi)調(diào)整。這也是這個(gè)方向值得投入的核心原因——適配層是數(shù)據(jù)管道里最常改動(dòng)、最需要快速迭代的部分把它獨(dú)立出來(lái)維護(hù)成本能降一個(gè)量級(jí)。3. 跑通第一條 pd2bs 轉(zhuǎn)換鏈路從配置到批量調(diào)用的最小腳本3.1 目錄結(jié)構(gòu)映射配置外置是第一原則我維護(hù)的 pd2bs-scripts 目錄結(jié)構(gòu)很樸素但每條規(guī)則都是踩過(guò)坑之后定下來(lái)的pd2bs-scripts/ ├── configs/ │ └── mappers.yaml # 字段映射配置改映射只動(dòng)這個(gè)文件 ├── input/ # 上游文件落地目錄 ├── bad/ # 校驗(yàn)失敗的數(shù)據(jù)與原因 ├── output/ # 轉(zhuǎn)換后的批次報(bào)文留作審計(jì) ├── logs/ # 運(yùn)行日志與批次統(tǒng)計(jì) ├── pd2bs.py # 主腳本 └── requirements.txtinput/ 目錄一般掛到上游管道同步路徑上上游文件到達(dá)后腳本即刻可見(jiàn)。output/ 目錄很多人覺(jué)得多余但它有兩個(gè)用處一是 BS 接口出問(wèn)題時(shí)不至于空口無(wú)憑直接把報(bào)文交給對(duì)方排查二是后面做對(duì)賬和回放時(shí)它是最可靠的事實(shí)記錄。映射配置外置是我最想強(qiáng)調(diào)的習(xí)慣。BS 接口的字段映射、枚舉轉(zhuǎn)換、默認(rèn)值全部放進(jìn) mappers.yaml一句話概括就是「改映射不改代碼」。這樣業(yè)務(wù)同事也能參與維護(hù)映射而不必每次找你改代碼。一個(gè)最小可用的 mappers.yaml 長(zhǎng)這樣# 映射規(guī)則target 是 BS 字段source 是 PD 字段 # type 可選string / int / amount_to_cents / datetime_utc / enum_map mappings: - target: outer_id source: order_id type: string - target: user_id source: user_id type: int - target: sku_code source: sku type: string - target: quantity source: num type: int - target: amount_cents source: amount type: amount_to_cents - target: paid_time source: paid_at type: datetime_utc timezone: Asia/Shanghai - target: status_code source: status type: enum_map enum_map: {PAID: 1, REFUNDED: 2, PENDING: 3} # 批次參數(shù) batch: size: 200 timeout: 30 max_retries: 3 base_delay: 0.5這里 timezone 指明上游時(shí)間的時(shí)區(qū)假設(shè)datetime_utc 處理器會(huì)按它解析再轉(zhuǎn) UTC。更重要的是枚舉映射沒(méi)有寫(xiě)在代碼里業(yè)務(wù)調(diào)整枚舉含義時(shí)只改配置即可。3.2 核心轉(zhuǎn)換讀文件、映射、類型轉(zhuǎn)換主腳本的核心是一個(gè)按配置逐字段轉(zhuǎn)換的函數(shù)。這里有一個(gè)關(guān)鍵設(shè)計(jì)用哨兵值標(biāo)記「字段缺失」而不是用 dict.get 默認(rèn)返回 None。區(qū)別我會(huì)在避坑章節(jié)細(xì)講先看代碼import json import yaml from datetime import datetime, timezone from zoneinfo import ZoneInfo _MISSING object() # 哨兵區(qū)分“字段缺失”和“字段值為 None” def load_mapping(path): with open(path, encodingutf-8) as f: cfg yaml.safe_load(f) return cfg[mappings], cfg[batch] def amount_to_cents(raw): # 元轉(zhuǎn)分用字符串運(yùn)算避免浮點(diǎn)誤差 return int(round(float(raw) * 100)) # 僅用于金額列確保 raw 是明確的數(shù)值字符串 def datetime_utc(raw, tz_name): if not raw: return None local datetime.strptime(raw, %Y-%m-%d %H:%M:%S) return local.replace(tzinfoZoneInfo(tz_name)).astimezone(timezone.utc).isoformat().replace(00:00, Z) def apply_mapping(row, mappings): out {} for rule in mappings: raw row.get(rule[source], _MISSING) if raw is _MISSING: out[rule[target]] None continue t rule.get(type, string) if t int: out[rule[target]] int(str(raw).strip()) elif t amount_to_cents: out[rule[target]] amount_to_cents(raw) elif t datetime_utc: out[rule[target]] datetime_utc(raw, rule.get(timezone, Asia/Shanghai)) elif t enum_map: out[rule[target]] rule[enum_map].get(raw) else: out[rule[target]] raw return out def load_jsonl(path): rows [] with open(path, encodingutf-8) as f: for line in f: line line.strip() if not line: continue rows.append(json.loads(line)) return rows這段代碼的邏輯很直白load_jsonl 按行讀入上游文件apply_mapping 對(duì)每一行執(zhí)行映射規(guī)則。值得說(shuō)明的是哨兵 _MISSING 的用法——row.get(source, _MISSING) 讓「字段不存在」和「字段值為 null」走不同分支。映射后值為 None 的字段在后續(xù)校驗(yàn)和發(fā)送環(huán)節(jié)會(huì)有專門(mén)處理而不是被默認(rèn)值悄悄替換掉。參數(shù)說(shuō)明type 決定轉(zhuǎn)換方式enum_map 里的字典可以隨時(shí)擴(kuò)展timezone 字段只在 datetime_utc 類型下生效。如果你的上游時(shí)間和時(shí)區(qū)假設(shè)變了只改配置不動(dòng)代碼。int 轉(zhuǎn)換前先 strip是為了對(duì)付上游偶爾出現(xiàn)的空格字符。3.3 校驗(yàn)與失敗兜底什么數(shù)據(jù)該攔在門(mén)外轉(zhuǎn)換完成不等于可以發(fā)送。BS 接口對(duì)數(shù)據(jù)的完整性校驗(yàn)很?chē)?yán)格與其讓接口返回一條錯(cuò)誤導(dǎo)致整批失敗不如在腳本側(cè)先攔住明顯有問(wèn)題的數(shù)據(jù)。我的校驗(yàn)函數(shù)只做四件事必填字段非空、數(shù)值范圍、枚舉合法、業(yè)務(wù)狀態(tài)檢查def validate_row(row): errors [] if not row.get(outer_id): errors.append(outer_id 為空) if row.get(quantity) is None or row.get(quantity) 0: errors.append(quantity 必須大于 0) if row.get(status_code) not in (1, 2, 3): errors.append(fstatus_code 非法: {row.get(status_code)}) if row.get(amount_cents) is None or row.get(amount_cents) 0: errors.append(amount_cents 非法) return errors def split_rows(rows, errors_map): good, bad [], [] for idx, row in enumerate(rows): errs validate_row(row) if errs: bad.append((idx, row, errs)) else: good.append(row) return good, bad校驗(yàn)規(guī)則本質(zhì)上是 BS 接口契約的本地切片。每一條規(guī)則都能對(duì)應(yīng)到接口文檔里的一句話比如「quantity 必須大于 0」對(duì)應(yīng)接口對(duì)訂購(gòu)數(shù)量的約束。這樣壞數(shù)據(jù)不會(huì)進(jìn)入網(wǎng)絡(luò)請(qǐng)求而是連同行號(hào)和原因一起寫(xiě)進(jìn) bad/ 目錄下的文件方便人工處理。失敗兜底我一般這樣寫(xiě)bad 文件命名帶上批次和日期內(nèi)容保留原始行和校驗(yàn)錯(cuò)誤列表。這樣上游拿到文件就能定位不用再跑一遍腳本看日志。這比把壞數(shù)據(jù)只打在 stdout 里靠譜得多。3.4 拼裝批量報(bào)文并調(diào)用分批的邊界條件轉(zhuǎn)換和校驗(yàn)之后就可以把數(shù)據(jù)送給 BS 了。編碼上要注意兩個(gè)細(xì)節(jié)用 Session 復(fù)用連接池分批大小從配置讀取而不是硬編碼import requests from requests.adapters import HTTPAdapter def build_payload(batch, batch_id, service_codeorder_sync): return { service_code: service_code, batch_id: batch_id, items: batch, } def send_batch(session, batch, batch_id, endpoint, timeout30): payload build_payload(batch, batch_id) resp session.post(endpoint, jsonpayload, timeouttimeout) resp.raise_for_status() return resp.json() def chunks(rows, size): for i in range(0, len(rows), size): yield rows[i:i size]調(diào)用方代碼就是把上面幾個(gè)函數(shù)串起來(lái)def main(input_path, cfg_path, endpoint): mappings, batch_cfg load_mapping(cfg_path) rows load_jsonl(input_path) mapped [apply_mapping(r, mappings) for r in rows] good, bad split_rows(mapped, {}) # 把 bad 寫(xiě)入 bad/ 目錄這里省略 session requests.Session() session.mount(endpoint, HTTPAdapter(max_retries0)) # 重試交給 call_with_retry for idx, batch in enumerate(chunks(good, batch_cfg[size])): batch_id fpd2bs_{input_path.stem}_{idx:03d} send_batch(session, batch, batch_id, endpoint, timeoutbatch_cfg[timeout])注意 HTTPAdapter 的 max_retries 我建議設(shè) 0把重試邏輯統(tǒng)一收口在應(yīng)用層這樣能精確控制退避策略和重試次數(shù)而不是依賴 requests 內(nèi)置的簡(jiǎn)單重試。batch_id 是冪等鍵的核心組成部分BS 側(cè)拿它做重復(fù)請(qǐng)求去重所以必須保證同一次轉(zhuǎn)換的每個(gè)批次都有唯一 ID重跑時(shí)也不能變。4. 調(diào)參實(shí)戰(zhàn)批量、并發(fā)、超時(shí)與重試怎么配才不翻車(chē)4.1 batch_size 不是越大越好接口超時(shí)與內(nèi)存的雙重約束第一批腳本上線時(shí)我天真地認(rèn)為 batch_size 越大越快直接配了接口文檔允許的上限 2000結(jié)果連續(xù)三批超時(shí)重試又疊加壓力BS 側(cè)告警響成一片。后來(lái)老老實(shí)實(shí)做了一組對(duì)比測(cè)試batch_size單批耗時(shí)p95現(xiàn)象501.2s請(qǐng)求數(shù)多總時(shí)長(zhǎng)被網(wǎng)絡(luò)往返稀釋2001.8s多數(shù)接口的甜點(diǎn)區(qū)間失敗重試成本可控8005.6s單批超時(shí)概率上升超時(shí)后整批重試代價(jià)高200012s內(nèi)存和序列化壓力大接口大概率 504結(jié)論很明確接口文檔說(shuō)的 max_items 是上限不是推薦值。我一般從接口允許值的一半起步用小批量樣本跑三組觀察 p95 耗時(shí)和錯(cuò)誤率再逐步往上加。同時(shí)要注意內(nèi)存batch_size 乘單條報(bào)文大小再乘并發(fā)數(shù)才是腳本的瞬時(shí)內(nèi)存峰值。200 條報(bào)文可能只有幾百 KB2000 條就可能到幾十 MB對(duì)常駐腳本來(lái)說(shuō)不算大但對(duì)跑批任務(wù)來(lái)說(shuō)沒(méi)必要冒這個(gè)風(fēng)險(xiǎn)。超時(shí)設(shè)置也要跟著 batch_size 走。batch 越大單批處理時(shí)間越長(zhǎng)timeout 不能還停留在 5 秒。我常用的經(jīng)驗(yàn)值timeout 設(shè)置為該批次正常耗時(shí)的 3 倍左右。比如 batch 200 正常 1.8 秒timeout 給 5 秒batch 800 正常 5.6 秒timeout 至少給 15 秒。timeout 太短會(huì)把慢請(qǐng)求誤判為失敗觸發(fā)無(wú)謂重試。4.2 重試策略指數(shù)退避、抖動(dòng)與冪等鍵缺一不可重試是轉(zhuǎn)換腳本最容易寫(xiě)壞的部分。常見(jiàn)做法是遇到任何異常都重試三次結(jié)果業(yè)務(wù)校驗(yàn)錯(cuò)誤被反復(fù)重試接口返回 400 還重試三次白白浪費(fèi)資源。我的原則是只有連接類異常和 5xx 才值得退避重試4xx 是客戶端問(wèn)題重試永遠(yuǎn)不會(huì)成功。import time import random def call_with_retry(fn, max_retries3, base_delay0.5): for attempt in range(max_retries 1): try: return fn() except ( requests.exceptions.ConnectTimeout, requests.exceptions.ConnectionError, requests.exceptions.HTTPError, ) as e: if attempt max_retries: raise # 5xx 和 429 由 HTTPError 拋出時(shí)按狀態(tài)碼區(qū)分 status getattr(e.response, status_code, None) if status is not None and status 500 and status ! 429: raise delay base_delay * (2 ** attempt) random.uniform(0, 0.2) time.sleep(delay) return None這里的指數(shù)退避是 0.5 秒、1 秒、2 秒遞增再加 0 到 0.2 秒的隨機(jī)抖動(dòng)。抖動(dòng)必須加否則多個(gè)并發(fā)批次同時(shí)失敗時(shí)重試也會(huì)同時(shí)發(fā)起形成另一種形式的驚群。429 特別說(shuō)明一下BS 返回 429 時(shí)通常帶 Retry-After 頭如果響應(yīng)里有這個(gè)字段應(yīng)該以它為準(zhǔn)而不是自己瞎猜等待時(shí)間。但所有重試的前提是冪等鍵。BS 接口必須支持按 batch_id 去重否則腳本重試一個(gè)已經(jīng)被部分處理的批次就會(huì)產(chǎn)生重復(fù)數(shù)據(jù)。對(duì)接 BS 時(shí)第一件事就要確認(rèn)接口是否冪等如果不支持腳本側(cè)就要在本地記錄已成功批次重跑前先查本地狀態(tài)。4.3 并發(fā)上限把腳本做成受控的消費(fèi)者而不是壓測(cè)工具跑批腳本很容易被人為加并發(fā)來(lái)提速但這個(gè)動(dòng)作要克制。BS 是業(yè)務(wù)系統(tǒng)它的容量不只是為你一個(gè)腳本準(zhǔn)備的同一時(shí)間可能還有別的任務(wù)在調(diào)用。我見(jiàn)過(guò)的一次事故就是轉(zhuǎn)換腳本開(kāi)了 16 個(gè)線程把 BS 的批量接口打到限流影響了線上正常業(yè)務(wù)。我常用的做法是先用單線程跑通確認(rèn)接口穩(wěn)了再用 ThreadPoolExecutor 逐步加并發(fā)。最大并發(fā)一般不超過(guò) 4而且要看 BS 側(cè)的容量評(píng)估。代碼上用一個(gè)信號(hào)量就能把整體并發(fā)封頂from concurrent.futures import ThreadPoolExecutor import threading sem threading.Semaphore(4) def bounded_send(batch, batch_id, session, endpoint, timeout): with sem: return send_batch(session, batch, batch_id, endpoint, timeout) with ThreadPoolExecutor(max_workers4) as pool: futures [ pool.submit(bounded_send, batch, batch_id, session, endpoint, timeout) for batch, batch_id in batches ] for f in futures: f.result()信號(hào)量和線程池的 max_workers 雙保險(xiǎn)主要防的是未來(lái)有人把 max_workers 改大時(shí)信號(hào)量還能兜住對(duì) BS 的最大并發(fā)。這種做法看著笨但跑批腳本的第一目標(biāo)是別惹麻煩而不是跑出性能壓測(cè)的架勢(shì)。數(shù)據(jù)量實(shí)在大的時(shí)候正確的方向是拆成多個(gè)窗口期任務(wù)而不是在一個(gè)腳本里無(wú)限堆并發(fā)。4.4 監(jiān)控日志與批次對(duì)賬腳本跑完看一眼退出碼是遠(yuǎn)遠(yuǎn)不夠的。批量轉(zhuǎn)換里最容易出現(xiàn)的問(wèn)題就是「整體成功個(gè)別失敗」而失敗記錄淹沒(méi)在日志里。我要求 pd2bs 每處理完一批就輸出一行結(jié)構(gòu)化日志2024-05-18 03:30:12 INFO batchpd2bs_20240518_001 items200 ok198 fail2 cost_ms1873 trace7f3a9c這一行的信息量很大items 是這批總量ok 和 fail 是 BS 返回的成功失敗數(shù)cost_ms 是耗時(shí)trace 是關(guān)聯(lián) ID。后續(xù)排查時(shí)按 trace 能找到 BS 側(cè)完整的處理鏈路按 batch 能找到本地 output/ 目錄留存的報(bào)文原文。fail 數(shù)不為 0 時(shí)腳本不應(yīng)該默默繼續(xù)。我習(xí)慣把失敗詳情單獨(dú)落一個(gè) CSV每行包括批次號(hào)、行號(hào)、業(yè)務(wù)主鍵、失敗原因方便上游和 BS 兩側(cè)一起定位。這個(gè) CSV 比對(duì)賬腳本還好用因?yàn)樗寝D(zhuǎn)換側(cè)和接口側(cè)事實(shí)的交叉點(diǎn)。沒(méi)有這批日志的跑批腳本出了事就是一個(gè)黑匣子只能靠猜。5. pd2bs 避坑實(shí)錄五個(gè)把轉(zhuǎn)換腳本搞掛的真實(shí)問(wèn)題這一章的內(nèi)容全是血淚經(jīng)驗(yàn)。每一條我都親自遇到過(guò)也跟著排過(guò)別人的類似問(wèn)題按「現(xiàn)象 → 原因 → 解決」寫(xiě)清楚。5.1 長(zhǎng)整型 ID 變成科學(xué)計(jì)數(shù)法float 轉(zhuǎn)換丟精度現(xiàn)象轉(zhuǎn)換后的 outer_id 在 BS 側(cè)查出來(lái)變成2.025e15這種樣子再轉(zhuǎn)回字符串就和原始值對(duì)不上了單號(hào)丟失最后幾位。原因上游的 order_id 是 19 位長(zhǎng)整型某個(gè)環(huán)節(jié)用了int()后又經(jīng)過(guò)一次 float 運(yùn)算或 JSON 序列化數(shù)字被轉(zhuǎn)成浮點(diǎn)浮點(diǎn)只能精確表示 2 的 53 次方以內(nèi)的整數(shù)超出部分直接丟精度。更隱蔽的路徑是 Excel 打開(kāi) CSV 時(shí)自動(dòng)轉(zhuǎn)成科學(xué)計(jì)數(shù)法再保存就不可逆了。解決所有 ID 字段全程按字符串處理。映射配置里 type 用 string不要用 int如果必須傳給 BS 整數(shù)型 ID先確認(rèn)位數(shù)在安全范圍內(nèi)并在轉(zhuǎn)換函數(shù)里加一個(gè)斷言if len(str(raw)) 15: raise ValueError。寧可腳本報(bào)錯(cuò)也不能讓錯(cuò)誤數(shù)據(jù)靜默流入下游。5.2 入庫(kù)時(shí)間差 8 小時(shí)本地時(shí)間與 UTC 的隱形邊界現(xiàn)象BS 側(cè)查到的 paid_time 普遍比實(shí)際支付時(shí)間晚或早了 8 小時(shí)但又不是所有行都差有的是 7 小時(shí)看著像隨機(jī)飄。原因上游導(dǎo)出的 paid_at 是2024-05-18 03:22:11沒(méi)有時(shí)區(qū)標(biāo)記。腳本里如果是用datetime.fromisoformat(raw).isoformat() Z直接拼等于把本地時(shí)間當(dāng)成了 UTC轉(zhuǎn)換后整體偏移。更鬧心的是如果 BS 側(cè)又做了一次解析時(shí)區(qū)判定不一致就會(huì)出現(xiàn) 7 小時(shí)、8 小時(shí)這種看似隨機(jī)的結(jié)果。解決解析時(shí)必須指定上游時(shí)區(qū)再轉(zhuǎn) UTC而不是直接拼字符local datetime.strptime(raw, %Y-%m-%d %H:%M:%S) utc local.replace(tzinfoZoneInfo(Asia/Shanghai)).astimezone(timezone.utc) result utc.strftime(%Y-%m-%dT%H:%M:%SZ)我踩過(guò)這個(gè)坑之后立了一條規(guī)矩所有時(shí)間字段必須在映射配置里顯式聲明 timezone腳本層禁止出現(xiàn)裸的 datetime 字符串拼接。上游換時(shí)區(qū)假設(shè)是配置變更而不是代碼變更。5.3 同一個(gè)文件兩種結(jié)果dict.get 和 or 混用的默認(rèn)值陷阱現(xiàn)象同樣的輸入文件跑兩次轉(zhuǎn)換一部分行的默認(rèn)值不一樣導(dǎo)致對(duì)賬不通過(guò)。排查半天發(fā)現(xiàn)是代碼分支不同。原因映射函數(shù)里有的地方寫(xiě)row.get(coupon, )有的地方寫(xiě)row.get(coupon) or 。當(dāng) coupon 字段存在但值為空字符串時(shí)前者保留空字符串后者把空字符串當(dāng)成假值替換成默認(rèn)值。如果還有row.get(num) or 0這種寫(xiě)法num 等于 0 的合法數(shù)據(jù)也會(huì)被替換成 0看似沒(méi)區(qū)別但 num 等于 None 和 num 等于 0 在 BS 側(cè)語(yǔ)義完全不同一個(gè)代表未填寫(xiě)一個(gè)代表真實(shí)數(shù)量。解決統(tǒng)一用哨兵 _MISSING 判斷「字段缺失」值和默認(rèn)值分清楚raw row.get(coupon, _MISSING) if raw is _MISSING: out[coupon] # 字段缺失時(shí)給默認(rèn)值 else: out[coupon] raw # 字段存在時(shí)原樣保留哪怕它是空串這條規(guī)則我寫(xiě)進(jìn)了代碼評(píng)審清單??吹給r出現(xiàn)在映射邏輯里基本都要打回去重寫(xiě)。5.4 空值把線上數(shù)據(jù)清空了更新語(yǔ)義下 None 不該出場(chǎng)現(xiàn)象某次同步后BS 側(cè)一批訂單的收貨地址變成空而原始數(shù)據(jù)里地址字段只是部分缺失不該覆蓋線上已有值。原因BS 的這個(gè)接口是「全量更新」語(yǔ)義報(bào)文字段缺省時(shí)接口默認(rèn)不更新但顯式傳 null 時(shí)接口會(huì)去更新該字段。pd2bs 腳本在字段缺失時(shí)映射為 None序列化 JSON 時(shí) null 被原樣帶出等于告訴 BS「把這幾個(gè)字段清空」。對(duì) insert 類接口這可能沒(méi)影響對(duì) update 類接口就是事故。解決把接口語(yǔ)義分成 insert 和 update 兩類update 場(chǎng)景下映射出的 None 字段在發(fā)送前剔除def strip_none_for_update(batch): cleaned [] for row in batch: item {k: v for k, v in row.items() if v is not None} cleaned.append(item) return cleaned同時(shí)在被剔除的字段里挑幾個(gè)業(yè)務(wù)關(guān)鍵字段記一條 WARN 日志。這樣既不影響更新語(yǔ)義也能在審計(jì)日志里留下線索知道哪些行哪些字段因?yàn)槿笔П惶^(guò)。5.5 重試風(fēng)暴腳本恢復(fù)后把下游打到限流現(xiàn)象腳本凌晨處理到一半掛了第二天補(bǔ)跑時(shí)所有失敗批次幾乎同時(shí)發(fā)起重試BS 接口直接限流連帶著正常業(yè)務(wù)請(qǐng)求也受影響。原因腳本掛掉時(shí)內(nèi)存里的所有批次狀態(tài)全部丟失。補(bǔ)跑邏輯如果簡(jiǎn)單粗暴地把全部批次重新投遞加上上一輪遺留的失敗重試疊加并發(fā)后瞬間打滿 BS。本質(zhì)是重試沒(méi)有全局限速每個(gè)批次各自為戰(zhàn)。解決補(bǔ)跑前先查本地 output/ 目錄和日志確認(rèn)哪些 batch_id 已經(jīng)成功只重跑失敗批次。同時(shí)加一層全局限速不管并發(fā)多少每秒最多發(fā)起固定數(shù)量的批次請(qǐng)求class RateLimiter: def __init__(self, max_per_second): self.min_interval 1.0 / max_per_second self.next_call 0 self.lock threading.Lock() def wait(self): with self.lock: now time.time() wait self.next_call - now if wait 0: time.sleep(wait) self.next_call max(self.next_call, time.time()) self.min_interval重試次數(shù)也壓到 2 次以內(nèi)超過(guò)就進(jìn)死信文件不再自動(dòng)重試。跑批腳本的生命在于可控寧可慢一點(diǎn)也不能因?yàn)樽约旱闹卦嚢严掠胃銙臁_@條是我在這個(gè)項(xiàng)目里交過(guò)最貴的一筆學(xué)費(fèi)。6. 進(jìn)階用法把 pd2bs 升級(jí)成可回放、可對(duì)賬的調(diào)度任務(wù)6.1 批次回放讓歷史文件可以原樣重跑跑批任務(wù)最怕的是「當(dāng)時(shí)跑過(guò)了但當(dāng)時(shí)的數(shù)據(jù)有問(wèn)題」。所以我在 input/ 文件處理完成后不刪除原文件只移動(dòng)到 input/archive/ 下按日期歸檔。腳本每次運(yùn)行都生成一個(gè)批次清單記錄輸入文件、輸出報(bào)文、批次號(hào)的對(duì)應(yīng)關(guān)系。重跑時(shí)直接用同樣參數(shù)再執(zhí)行一遍由于 batch_id 由文件名和序號(hào)生成重跑結(jié)果和第一次完全一致BS 側(cè)靠?jī)绲孺I自動(dòng)忽略重復(fù)數(shù)據(jù)。這個(gè)設(shè)計(jì)給排查問(wèn)題提供了后悔藥。某次 BS 側(cè)數(shù)據(jù)異常懷疑是轉(zhuǎn)換邏輯寫(xiě)錯(cuò)我只要把當(dāng)時(shí)的映射配置和輸入文件都翻出來(lái)重跑一次對(duì)比 output/ 里的報(bào)文就能確認(rèn)是腳本問(wèn)題還是接口問(wèn)題。沒(méi)有回放能力遇到這種問(wèn)題就只能靠嘴對(duì)線。6.2 對(duì)賬命令轉(zhuǎn)換正確性的最后防線我習(xí)慣在 pd2bs 腳本里加一個(gè)--reconcile模式只做統(tǒng)計(jì)不調(diào)用接口。它把 input/ 和 output/ 各算一遍總條數(shù)、總金額、狀態(tài)分布然后對(duì)比。一行命令就能看出轉(zhuǎn)換環(huán)節(jié)有沒(méi)有丟數(shù)據(jù)python pd2bs.py --reconcile --input input/20240518_orders.jsonl --output output/對(duì)賬結(jié)果會(huì)輸出一個(gè)三行的小表原始行數(shù)、轉(zhuǎn)換后行數(shù)、失敗行數(shù)。金額合計(jì)從元轉(zhuǎn)換成分之后應(yīng)當(dāng)完全相等。這個(gè)動(dòng)作建議每次批次跑完后自動(dòng)執(zhí)行一次連續(xù)兩天對(duì)不上賬說(shuō)明有靜默丟失早點(diǎn)暴露比下游投訴時(shí)才發(fā)現(xiàn)要好得多。6.3 一個(gè)讓我長(zhǎng)記性的習(xí)慣我吃過(guò)一次教訓(xùn)某次改映射配置把枚舉值 PAID 的映射數(shù)字寫(xiě)錯(cuò)結(jié)果整批訂單的 status_code 全部變成另一個(gè)狀態(tài)。當(dāng)時(shí)沒(méi)有做全量對(duì)賬只看了腳本退出碼為 0 就放它跑了等到業(yè)務(wù)側(cè)發(fā)現(xiàn)異常已經(jīng)過(guò)去了大半天。從那以后我立了一個(gè)習(xí)慣任何映射配置改動(dòng)先用最近一天的輸入文件跑一遍小樣本對(duì)賬確認(rèn)枚舉、金額、時(shí)間三類字段的分布與預(yù)期一致再跑全量。這個(gè)動(dòng)作成本很低但能攔住絕大多數(shù)映射層面的低級(jí)錯(cuò)誤。pd2bs 這類腳本的價(jià)值不在于代碼寫(xiě)得多漂亮而在于它讓數(shù)據(jù)管道下游變得可控、可查、可重來(lái)。每次調(diào)整映射時(shí)多問(wèn)一句「這次改動(dòng)影響哪些字段」每次跑批后多看一眼對(duì)賬統(tǒng)計(jì)累積下來(lái)省下的排查時(shí)間遠(yuǎn)超寫(xiě)腳本的時(shí)間。希望幫到你。本文還有配套的精品資源點(diǎn)擊獲取