消息路由與任務(wù)執(zhí)行引擎 hermes-agent 的設(shè)計(jì)與實(shí)踐)
1. 項(xiàng)目緣起為什么我會(huì)自己寫(xiě)一個(gè)叫 hermes-agent 的東西先聊聊名字。Hermes 是希臘神話里的信使神腳上長(zhǎng)著翅膀負(fù)責(zé)在眾神之間傳遞消息。我當(dāng)時(shí)給這個(gè)項(xiàng)目起名的時(shí)候想的就是這一層含義——在一個(gè)系統(tǒng)越來(lái)越復(fù)雜的時(shí)代消息的傳遞、路由、分發(fā)、執(zhí)行就是數(shù)字世界的信使工作。所以我管它叫 hermes-agent一個(gè)專門負(fù)責(zé)消息接入、規(guī)則匹配和自動(dòng)執(zhí)行的小型智能體框架。這個(gè)項(xiàng)目的背景很具體。我手里維護(hù)著好幾套內(nèi)部服務(wù)有定時(shí)任務(wù)、有外部回調(diào)、有用戶行為事件還有各類監(jiān)控告警。這些消息源格式五花八門有的是 JSON有的帶簽名需要驗(yàn)簽有的是純文本塞在 Webhook 里。早期我的處理方式是在每個(gè)服務(wù)里各寫(xiě)各的入口邏輯A 服務(wù)收到回調(diào)就更新數(shù)據(jù)庫(kù)B 服務(wù)收到事件就調(diào)用某個(gè)內(nèi)部 APIC 服務(wù)收到告警就發(fā)釘釘通知。問(wèn)題是一旦消息處理邏輯發(fā)生變化我需要同時(shí)改好幾個(gè)服務(wù)重新部署排錯(cuò)的時(shí)候還要挨個(gè)翻日志非常被動(dòng)。后來(lái)我認(rèn)真想了一下這些場(chǎng)景的本質(zhì)其實(shí)是一樣的源頭把事件拋出來(lái)中間需要有一套穩(wěn)定的機(jī)制接住再根據(jù)規(guī)則決定接下來(lái)讓誰(shuí)去干活。這正是 agent 最擅長(zhǎng)的事情。于是 hermes-agent 的雛形誕生了——一個(gè)運(yùn)行時(shí)基于 Python asyncio、存儲(chǔ)用 SQLite生產(chǎn)環(huán)境可以換 Postgres、規(guī)則配置走 YAML 的輕量級(jí)消息代理與任務(wù)執(zhí)行引擎。它不依賴重量級(jí)框架部署就是拉代碼起進(jìn)程非常適合中小型團(tuán)隊(duì)內(nèi)部自動(dòng)化場(chǎng)景。這篇文章我會(huì)把整個(gè)項(xiàng)目的設(shè)計(jì)思路、運(yùn)行模型、核心代碼、踩坑記錄和演進(jìn)方向都攤開(kāi)來(lái)講。如果你也在做 Webhook 統(tǒng)一接入、任務(wù)分發(fā)、事件驅(qū)動(dòng)型自動(dòng)化或者單純想看看一個(gè) agent 框架應(yīng)該怎么設(shè)計(jì)這篇內(nèi)容應(yīng)該對(duì)你有幫助。對(duì)于基礎(chǔ)偏弱的讀者我會(huì)盡量把涉及到的基礎(chǔ)概念一并講清楚不預(yù)設(shè)你已經(jīng)掌握 Python 異步編程或者消息隊(duì)列的完整知識(shí)。2. 核心架構(gòu)消息路由與任務(wù)編排的運(yùn)行模型2.1 三層模塊劃分接入層、規(guī)則層、執(zhí)行層hermes-agent 從設(shè)計(jì)之初就是三層結(jié)構(gòu)這個(gè)分層思路不是我拍腦袋決定的而是從實(shí)際需求里自然長(zhǎng)出來(lái)的。第一層是接入層Collector。它負(fù)責(zé)統(tǒng)一接住各種來(lái)源的消息。常見(jiàn)的 Collector 有 HTTP Webhook Collector、定時(shí)器 Collector、文件監(jiān)聽(tīng) Collector 和隊(duì)列消費(fèi) Collector。對(duì)于 HTTP 類型的 Collector我會(huì)封裝成一個(gè)獨(dú)立的 Web 服務(wù)進(jìn)程只做一件事把收到的請(qǐng)求體轉(zhuǎn)成統(tǒng)一格式的 Event 對(duì)象寫(xiě)入消息隊(duì)列或存儲(chǔ)然后立刻返回響應(yīng)給上游。這樣上游接口的響應(yīng)速度不會(huì)受到下游處理邏輯的影響這是個(gè)很關(guān)鍵的設(shè)計(jì)。第二層是規(guī)則層Router。它負(fù)責(zé)判斷一條消息應(yīng)該交給誰(shuí)處理。判斷依據(jù)可以從消息內(nèi)容、消息來(lái)源、消息優(yōu)先級(jí)等多個(gè)維度提取。每一條規(guī)則本質(zhì)上是一組條件和動(dòng)作的映射。條件支持等值匹配、正則匹配、JSONPath 取值匹配等多種模式。動(dòng)作則是指定執(zhí)行器名稱和參數(shù)。第三層是執(zhí)行層Executor。它負(fù)責(zé)真正干活。執(zhí)行器是一段可插拔的代碼單元接收 Event 對(duì)象執(zhí)行具體業(yè)務(wù)邏輯。比如 HttpExecutor 把 Event 轉(zhuǎn)成 HTTP 請(qǐng)求發(fā)給內(nèi)部系統(tǒng)PythonExecutor 直接執(zhí)行一段配置好的 Python 腳本ShellExecutor 跑一條命令NotifyExecutor 發(fā)送通知到釘釘或者郵件。這樣的好處是接新場(chǎng)景時(shí)你不需要?jiǎng)涌蚣鼙旧碇恍枰獙?xiě)一個(gè)新的 Executor 注冊(cè)進(jìn)去再配一條規(guī)則。這三層之間通過(guò)一個(gè)消息存儲(chǔ)默認(rèn) SQLite生產(chǎn)建議 Postgres和進(jìn)程內(nèi)的消息通道串起來(lái)。實(shí)際運(yùn)行的時(shí)候接入層進(jìn)程把消息寫(xiě)進(jìn)存儲(chǔ)并標(biāo)記為 pending處理進(jìn)程輪詢 pending 消息交給規(guī)則層判斷再調(diào)用執(zhí)行層完成任務(wù)并回寫(xiě)狀態(tài)。2.2 一次消息從接入到執(zhí)行完成的完整生命周期我用一個(gè)實(shí)際例子說(shuō)明整個(gè)流程。假設(shè)你有一個(gè)電商訂單系統(tǒng)用戶支付成功后訂單服務(wù)會(huì)調(diào)用你配置的 Webhook 地址通知支付結(jié)果。這個(gè)地址就是 hermes-agent 暴露出的 HTTP 接入點(diǎn)。第一步訂單服務(wù) POST 一條 JSON 到/hooks/payment內(nèi)容大概是{order_id: ORD20250101, status: paid, amount: 99.5}。HTTP Collector 收到請(qǐng)求后先做格式校驗(yàn)提取必要的元信息來(lái)源標(biāo)記、接收時(shí)間、原始內(nèi)容然后封裝成 Event 對(duì)象dataclass class Event: event_id: str # 全局唯一ID用于冪等 source: str # 來(lái)源標(biāo)記比如 payment_webhook event_type: str # 事件類型比如 payment.paid payload: dict # 原始消息內(nèi)容 received_at: datetime priority: int # 優(yōu)先級(jí)默認(rèn)0越大越優(yōu)先第二步Collector 把 Event 寫(xiě)入消息存儲(chǔ)。同時(shí)進(jìn)程內(nèi)會(huì)立即觸發(fā)一次路由判斷但注意這里不是同步阻塞的。寫(xiě)入成功后 HTTP 接口立刻返回 200處理邏輯全部在后臺(tái)繼續(xù)。這個(gè)異步設(shè)計(jì)能保證上游的 Webhook 調(diào)用不會(huì)被你這邊慢邏輯拖死。第三步處理進(jìn)程從存儲(chǔ)中輪詢到這條 pending 消息進(jìn)入 Router。Router 讀取已加載的規(guī)則集合并找到匹配項(xiàng)。rules: - name: payment_paid_order_finish match: event_type: payment.paid actions: - executor: python_executor params: script_path: ./scripts/order_finish.py timeout: 30第四步執(zhí)行器運(yùn)行后腳本讀取 Event把訂單狀態(tài)更新為已完成、發(fā)送短信通知用戶、調(diào)用倉(cāng)儲(chǔ)系統(tǒng)的出庫(kù)接口。執(zhí)行結(jié)果回寫(xiě)到消息存儲(chǔ)狀態(tài)從 pending 變成 success 或者 failed。如果失敗會(huì)進(jìn)入重試隊(duì)列按配置的退避策略重新執(zhí)行。整個(gè)生命周期最核心的原則是消息不丟失。消息無(wú)論是在接收階段還是在處理階段每一步的狀態(tài)更新都落盤保存絕不在內(nèi)存里直接處理完不記錄。這樣即使進(jìn)程中途崩潰重啟之后依然可以從 failed 或者 pending 狀態(tài)恢復(fù)執(zhí)行。2.3 數(shù)據(jù)模型與存儲(chǔ)選型邏輯存儲(chǔ)這塊我是從簡(jiǎn)單出發(fā)最初直接用 SQLite。直到現(xiàn)在如果場(chǎng)景是單機(jī)處理SQLite 完全夠用。表結(jié)構(gòu)非常簡(jiǎn)潔CREATE TABLE events ( id INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT UNIQUE NOT NULL, source TEXT NOT NULL, event_type TEXT NOT NULL, payload TEXT NOT NULL, status TEXT NOT NULL DEFAULT pending, priority INTEGER NOT NULL DEFAULT 0, retry_count INTEGER NOT NULL DEFAULT 0, created_at TIMESTAMP NOT NULL, updated_at TIMESTAMP NOT NULL ); CREATE INDEX idx_events_status_priority ON events(status, priority, created_at);這個(gè)表承擔(dān)了消息隊(duì)列和狀態(tài)存儲(chǔ)的雙重職責(zé)。很多剛接觸這個(gè)項(xiàng)目的朋友會(huì)問(wèn)為什么不直接用 RabbitMQ 或者 Kafka答案是看場(chǎng)景。對(duì)于內(nèi)部自動(dòng)化、Webhook 統(tǒng)一接入這類日吞吐量在幾萬(wàn)條以下的任務(wù)引入獨(dú)立消息中間件會(huì)增加部署成本和運(yùn)維成本。而且 SQLite 表天然支持 SQL 查詢消息審計(jì)、問(wèn)題排查非常方便。當(dāng)你真的需要橫向擴(kuò)容時(shí)把這個(gè)表?yè)Q到 Postgres再對(duì)處理進(jìn)程做多實(shí)例部署改動(dòng)成本并不高。選擇存儲(chǔ)方案的一個(gè)重要標(biāo)準(zhǔn)是你的消息量級(jí)是否值得引入一套額外的分布式系統(tǒng)。如果答案是不確定就用最簡(jiǎn)單可靠的方案起步。這個(gè)原則后來(lái)幫我省了非常多的時(shí)間。3. 規(guī)則引擎核心中的核心3.1 規(guī)則配置與條件匹配的設(shè)計(jì)思路Router 是整個(gè) hermes-agent 的決策大腦而規(guī)則引擎中的匹配條件則是大腦的每一個(gè)神經(jīng)元。我把條件匹配設(shè)計(jì)成三種形式覆蓋了絕大多數(shù)實(shí)際需求。第一種是 exact 等值匹配。就是消息某個(gè)字段的值等于某個(gè)固定值。適合場(chǎng)景來(lái)源是某個(gè)固定系統(tǒng)、事件類型是某個(gè)固定字符串。第二種是 regex 正則匹配。適合場(chǎng)景訂單號(hào)以 AB 開(kāi)頭、用戶郵箱屬于某個(gè)域名、請(qǐng)求路徑符合某個(gè)模式。正則匹配可以配置在 JSONPath 取值后的字符串上。第三種是 expression 表達(dá)式匹配。這是最靈活但也是我最謹(jǐn)慎使用的一種。你可以寫(xiě)一段簡(jiǎn)單的布爾表達(dá)式內(nèi)部通過(guò)受限的 eval 執(zhí)行不提供任意代碼執(zhí)行能力。屬性取值通過(guò) JSONPath 實(shí)現(xiàn)。統(tǒng)一配置方式如下rules: - name: high_value_order match: $type: all conditions: - field: $.payload.amount op: gt value: 5000 - field: $.event_type op: eq value: payment.paid priority: 90 actions: - executor: notify_executor params: channel: dingtalk template: high_value_order這條規(guī)則的意思是當(dāng)消息的 payload.amount 大于 5000且 event_type 等于 payment.paid就觸發(fā)一個(gè)高價(jià)值訂單的通知?jiǎng)幼?。priority 決定多條規(guī)則同時(shí)匹配時(shí)誰(shuí)先執(zhí)行。這個(gè)字段在設(shè)計(jì)早期被忽略過(guò)后來(lái)有人反饋說(shuō)同時(shí)命中財(cái)務(wù)通知和庫(kù)存通知時(shí)希望財(cái)務(wù)先跑才補(bǔ)上的。3.2 規(guī)則匹配的反向查詢優(yōu)化規(guī)則多了以后如果每條消息都遍歷所有規(guī)則去匹配性能會(huì)逐漸變差。我后來(lái)做了一項(xiàng)優(yōu)化反向索引。具體做法是啟動(dòng)時(shí)把所有規(guī)則的 event_type 條件提取出來(lái)建立字典然后按 event_type 快速定向到候選規(guī)則再進(jìn)一步做深度匹配。這個(gè)優(yōu)化把規(guī)則匹配的復(fù)雜度從 O(N) 降到接近 O(1)。實(shí)現(xiàn)上也很簡(jiǎn)單就是預(yù)處理階段class RuleIndex: def __init__(self, rules): self.event_type_map {} for rule in rules: for each_type in rule.get_match_event_types(): self.event_type_map.setdefault(each_type, []).append(rule) def find_candidates(self, event): types get_event_types(event.event_type) rules [] for t in types: rules.extend(self.event_type_map.get(t, [])) return deduplicate(rules)對(duì)于一條消息先用 event_type 快速篩出候選規(guī)則集合然后再逐條深度檢查完整條件極大減少了不必要的規(guī)則遍歷。從實(shí)際效果看規(guī)則數(shù)在三位數(shù)以內(nèi)時(shí)匹配耗時(shí)基本可以忽略。3.3 匹配失敗與未命中消息的處理策略做過(guò)通知系統(tǒng)的朋友應(yīng)該都有經(jīng)驗(yàn)最怕的就是消息沒(méi)匹配到規(guī)則然后靜默丟失。你根本不知道它丟了直到業(yè)務(wù)方來(lái)問(wèn)為什么我沒(méi)有收到消息。hermes-agent 對(duì)這一塊的處理是未命中的消息不丟棄統(tǒng)一進(jìn)入 unmatched 表同時(shí)提供管理接口可以查詢。這樣你可以定期復(fù)盤是否有新的消息類型需要添加規(guī)則。我見(jiàn)過(guò)不少系統(tǒng)從一開(kāi)始的忽略異常到后來(lái)接二連三出問(wèn)題根源就是沒(méi)有負(fù)面消息的可見(jiàn)性。我的建議是把未命中消息當(dāng)成一等公民對(duì)待它們和正常消息一樣重要。4. 執(zhí)行器的設(shè)計(jì)與實(shí)踐從腳本到插件的演進(jìn)4.1 執(zhí)行器接口設(shè)計(jì)與注冊(cè)管理執(zhí)行器是 hermes-agent 真正干活的部分。我定義了一個(gè)非常簡(jiǎn)潔的抽象接口class BaseExecutor: async def execute(self, event: Event, params: dict) - ExecResult: raise NotImplementedError每個(gè)執(zhí)行器只需要實(shí)現(xiàn)一個(gè) execute 方法。事件通過(guò)參數(shù)傳入執(zhí)行參數(shù)通過(guò) params 傳入。返回值是一個(gè) ExecResult包含 success、output、message、duration_ms 等字段方便持久化和追蹤。注冊(cè)機(jī)制方面我通過(guò)一個(gè)全局注冊(cè)表實(shí)現(xiàn)EXECUTOR_REGISTRY {} def register_executor(name): def decorator(cls): EXECUTOR_REGISTRY[name] cls return cls return decorator然后在每個(gè)執(zhí)行器文件里加上裝飾器??蚣軉?dòng)時(shí)把所有執(zhí)行器文件 import 一遍注冊(cè)表里就有了全量執(zhí)行器。規(guī)則配置里指定 executor 名稱Router 就能根據(jù)注冊(cè)表找到對(duì)應(yīng)類并實(shí)例化調(diào)用。新增執(zhí)行器完全不需要修改框架代碼做到了可插拔。4.2 幾個(gè)常用執(zhí)行器的實(shí)現(xiàn)細(xì)節(jié)HttpExecutor 是最常用的執(zhí)行器之一。它負(fù)責(zé)把事件轉(zhuǎn)發(fā)給其他內(nèi)部服務(wù)。具體實(shí)現(xiàn)時(shí)有兩個(gè)細(xì)節(jié)容易踩坑超時(shí)控制和重試策略。register_executor(http_executor) class HttpExecutor(BaseExecutor): async def execute(self, event, params): url params.get(url) method params.get(method, POST) payload event.payload timeout params.get(timeout, 10) async with aiohttp.ClientSession() as session: try: async with session.request(method, url, jsonpayload, timeoutaiohttp.ClientTimeout(totaltimeout)) as resp: body await resp.text() return ExecResult(successTrue, outputbody, messagefstatus{resp.status}) except asyncio.TimeoutError: return ExecResult(successFalse, messageftimeout after {timeout}s)重試不用在 HttpExecutor 里做因?yàn)橹卦囀窍⑻幚砜蚣軐用娴穆氊?zé)。執(zhí)行器只負(fù)責(zé)返回結(jié)果框架根據(jù)結(jié)果和配置決定要不要重試。這個(gè)職責(zé)邊界從設(shè)計(jì)之初就明確下來(lái)了避免執(zhí)行器里塞進(jìn)太多與業(yè)務(wù)無(wú)關(guān)的邏輯。PythonExecutor 比較特殊它允許在規(guī)則配置里指定一段腳本路徑框架動(dòng)態(tài)加載并執(zhí)行。這里有一個(gè)安全邊界腳本等同于本地代碼執(zhí)行權(quán)限等同于框架進(jìn)程本身。所以它只適合內(nèi)部環(huán)境下執(zhí)行可信代碼不適合對(duì)外開(kāi)放成多租戶執(zhí)行能力。NotifyExecutor 負(fù)責(zé)發(fā)送通知。支持釘釘、企業(yè)微信、郵件、飛書(shū)這些常見(jiàn)的通知渠道??紤]到外部 API 不穩(wěn)定這類執(zhí)行器最容易失敗所以通知類消息的重試策略通常配置得比較積極。4.3 執(zhí)行器的超時(shí)控制、并發(fā)與資源限制處理消息的時(shí)候最怕某個(gè)執(zhí)行器出問(wèn)題導(dǎo)致整個(gè)進(jìn)程卡死。所以從第一天起超時(shí)控制就是執(zhí)行器運(yùn)行的關(guān)鍵約束??蚣転槊看螆?zhí)行器調(diào)用都包裹了 asyncio.wait_for強(qiáng)制超時(shí)try: result await asyncio.wait_for(executor_instance.execute(event, params), timeouttimeout) except asyncio.TimeoutError: result ExecResult(successFalse, messageexecutor timeout)并發(fā)控制方面hermes-agent 支持配置最大并發(fā)數(shù)。默認(rèn)是 20也就是同一時(shí)間最多 20 個(gè)任務(wù)在跑。超過(guò)的部分排隊(duì)等待。這個(gè)設(shè)計(jì)是為了避免某個(gè)瞬間大量消息涌入時(shí)打爆下游系統(tǒng)。對(duì)于危險(xiǎn)操作類的執(zhí)行器比如刪除文件、清理數(shù)據(jù)庫(kù)數(shù)據(jù)我建議在 Executor 內(nèi)部增加二次確認(rèn)邏輯。規(guī)則配置里需要顯式設(shè)置confirm: true否則直接返回失敗。雖然增加了配置復(fù)雜度但在生產(chǎn)環(huán)境里這個(gè)保障非常重要。5. 可靠性設(shè)計(jì)重試、冪等與死信隊(duì)列5.1 重試策略與退避算法消息處理不可能永遠(yuǎn)一次成功。網(wǎng)絡(luò)抖動(dòng)、下游服務(wù)重啟、第三方接口超時(shí)這些都是常態(tài)。重試策略是整個(gè)可靠性設(shè)計(jì)中最重要的部分。hermes-agent 里的重試策略是四級(jí)退避第一次失敗后等 5 秒第二次 30 秒第三次 5 分鐘第四次 30 分鐘。超過(guò)四次仍然失敗消息進(jìn)入死信隊(duì)列。這個(gè)策略的經(jīng)驗(yàn)依據(jù)是大部分瞬時(shí)故障在前幾次重試時(shí)就能恢復(fù)如果 30 分鐘后還不行大概率不是瞬時(shí)問(wèn)題了沒(méi)必要無(wú)限重試。實(shí)現(xiàn)上我把重試信息放在消息的元數(shù)據(jù)里dataclass class RetryPolicy: max_retries: int 4 base_delay: int 5 multiplier: int 6 max_delay: int 1800每次重試的延遲時(shí)間 base_delay * multiplier^retry_count不超過(guò) max_delay。這個(gè)算法簡(jiǎn)單明了不需要引入復(fù)雜的指數(shù)退避庫(kù)。5.2 冪等設(shè)計(jì)避免重復(fù)執(zhí)行帶來(lái)副作用重試機(jī)制帶來(lái)的一個(gè)直接問(wèn)題是消息可能被處理不止一次。比如執(zhí)行器成功了但是由于網(wǎng)絡(luò)原因響應(yīng)沒(méi)有及時(shí)寫(xiě)回框架判斷超時(shí)后重試于是同一個(gè)事件被執(zhí)行了兩次。這是所有分布式任務(wù)系統(tǒng)都躲不開(kāi)的問(wèn)題唯一的解法就是冪等設(shè)計(jì)。冪等有兩種做法。第一種是業(yè)務(wù)層冪等你在自己的業(yè)務(wù)代碼里判斷這個(gè)訂單是否已經(jīng)被處理過(guò)了處理過(guò)就什么都不做直接返回成功。這種做法需要業(yè)務(wù)方配合。第二種是框架層冪等框架記錄每條消息的執(zhí)行指紋指紋相同的結(jié)果直接用不實(shí)際執(zhí)行。hermes-agent 在框架層面做了一個(gè)簡(jiǎn)單的保護(hù)。給每條事件生成 event_id這個(gè) ID 在消息源頭生成并隨消息攜帶。事件表里 event_id 是唯一索引同一事件重復(fù)寫(xiě)入直接失敗。處理成功之后這個(gè) ID 會(huì)回寫(xiě)到一個(gè)已處理表再次收到同 ID 的事件時(shí)直接返回成功。但要注意框架層的冪等保護(hù)沒(méi)法覆蓋執(zhí)行器已經(jīng)開(kāi)始干活但是結(jié)果沒(méi)記錄這種情況。所以我的建議是所有 Executor 在編寫(xiě)時(shí)都把事件已經(jīng)處理過(guò)當(dāng)成正常情況來(lái)處理不要拋出異常。把冪等當(dāng)成一種設(shè)計(jì)習(xí)慣而不是框架約束。5.3 死信隊(duì)列與人工介入流程死信隊(duì)列DLQ是消息系統(tǒng)的最后一道防線。我已經(jīng)設(shè)置了合理的重試次數(shù)和退避策略如果消息最終還是處理失敗就說(shuō)明靠自動(dòng)重試解決不了問(wèn)題。此時(shí)消息進(jìn)入 DLQ等待人工介入或后續(xù)補(bǔ)償腳本。在 hermes-agent 的管理界面上DLQ 消息可見(jiàn)、可查詢、可重新投遞。我通常的做法是先查詢 DLQ 里的消息內(nèi)容和失敗原因確認(rèn)問(wèn)題后修復(fù)比如配置錯(cuò)誤、依賴服務(wù)恢復(fù)然后在界面上點(diǎn)擊重新投遞消息會(huì)重新進(jìn)入 pending 狀態(tài)再走一遍處理流程。這里我給一個(gè)建議DLQ 里的消息需要定期檢查不要讓它靜默增長(zhǎng)??梢栽?hermes-agent 旁邊掛一個(gè)簡(jiǎn)單的 cron每天統(tǒng)計(jì) DLQ 數(shù)量超過(guò)閾值就發(fā)告警給你。不然時(shí)間久了積壓的失敗消息會(huì)變成一座沒(méi)人想動(dòng)的山。6. 實(shí)戰(zhàn)項(xiàng)目搭建一個(gè)統(tǒng)一的 Webhook 接收器6.1 場(chǎng)景與需求定義這個(gè)實(shí)戰(zhàn)項(xiàng)目來(lái)自一個(gè)真實(shí)需求。我們內(nèi)部有好幾個(gè)服務(wù)需要暴露 Webhook 給第三方調(diào)用包括支付回調(diào)、物流狀態(tài)回傳、短信狀態(tài)報(bào)告。每個(gè)服務(wù)的回調(diào)地址不同、驗(yàn)簽方式不同、處理邏輯不同。最麻煩的是某些第三方平臺(tái)的 Webhook 對(duì)響應(yīng)時(shí)間要求極高超過(guò) 2 秒就判定失敗并開(kāi)始重試。需求匯總下來(lái)有三點(diǎn)統(tǒng)一 Webhook 入口一個(gè)地址接收所有第三方事件。回調(diào)立刻返回 200具體邏輯異步處理。驗(yàn)簽邏輯可配置消息處理結(jié)果可追溯。6.2 配置實(shí)戰(zhàn)多來(lái)源接入與驗(yàn)簽我在 hermes-agent 里配置了三個(gè) Collector分別監(jiān)聽(tīng)三個(gè)路徑/hooks/payment、/hooks/logistics、/hooks/sms。每個(gè) Collector 綁定一個(gè)來(lái)源名稱同時(shí)可以配置該來(lái)源的驗(yàn)簽規(guī)則。為了演示我用支付回調(diào)的 HMAC-SHA256 簽名方式作為例子。collectors: - name: payment_hook path: /hooks/payment source: payment_callback verify: type: hmac_sha256 secret_env: PAYMENT_HOOK_SECRET header: X-Signature body_from: raw驗(yàn)簽邏輯不算復(fù)雜。第三方平臺(tái)用密鑰對(duì)請(qǐng)求體做 HMAC-SHA256 簽名放在 Header 里傳過(guò)來(lái)。hermes-agent 收到請(qǐng)求后用同樣的算法和密鑰計(jì)算簽名和 Header 里的值做比對(duì)。不一致就返回 401保持一致則繼續(xù)處理。這個(gè)驗(yàn)簽過(guò)程看起來(lái)簡(jiǎn)單但實(shí)際有一個(gè)非常重要的細(xì)節(jié)驗(yàn)簽必須使用原始請(qǐng)求體不能先解碼 JSON 再重新序列化。因?yàn)?JSON 鍵值順序變化會(huì)導(dǎo)致簽名不一致。我在代碼里專門保留了原始 body 字節(jié)用于驗(yàn)簽解析 JSON 放在驗(yàn)簽通過(guò)之后。這個(gè)坑非常隱蔽希望看到這里的讀者能記住。6.3 實(shí)戰(zhàn)踩坑第三方回調(diào)的響應(yīng)超時(shí)問(wèn)題這個(gè)項(xiàng)目的第一個(gè)線上問(wèn)題就出在回調(diào)響應(yīng)超時(shí)上。某個(gè)物流平臺(tái)的 Webhook 配置了 HTTP 超時(shí)時(shí)間為 2 秒。hermes-agent 收到回調(diào)后需要驗(yàn)簽、解析、寫(xiě)入 SQLite 數(shù)據(jù)庫(kù)、再返回 200。正常情況下整個(gè)過(guò)程在幾十毫秒內(nèi)完成。但某個(gè)時(shí)間段內(nèi)由于服務(wù)器的磁盤 I/O 出現(xiàn)偶發(fā)高延遲SQLite 寫(xiě)入超過(guò)了 2 秒導(dǎo)致上游平臺(tái)判定超時(shí)并重試。重試又帶來(lái)了重復(fù)消息處理邏輯里冪等沒(méi)做好部分物流狀態(tài)被重復(fù)回寫(xiě)。這個(gè)問(wèn)題的教訓(xùn)是雙重的。第一Webhook 接口的數(shù)據(jù)庫(kù)寫(xiě)入不能放在返回響應(yīng)的路徑上。我的修復(fù)方案是HTTP Collector 先寫(xiě)內(nèi)存隊(duì)列立刻返回 200后臺(tái)異步批量沖刷到 SQLite。內(nèi)存隊(duì)列的可靠性可以通過(guò)定時(shí)持久化快照來(lái)保障。第二冪等保護(hù)在任何可能重試的入口都必須做好即使你認(rèn)為概率很低。7. 性能優(yōu)化與壓測(cè)結(jié)果7.1 基于 asyncio 的并發(fā)模型背后的選擇理由我在最初選擇 asyncio 而不是多線程核心原因是這個(gè)項(xiàng)目的 IO 密集特性非常明顯接收 HTTP 請(qǐng)求、讀寫(xiě) SQLite、調(diào)用第三方接口、發(fā)送通知。這些都是 IO 等待沒(méi)有明顯的 CPU 密集計(jì)算。asyncio 在單線程內(nèi)用事件循環(huán)處理海量并發(fā)開(kāi)銷遠(yuǎn)小于線程切換且沒(méi)有線程安全問(wèn)題。代價(jià)是要小心不要寫(xiě)阻塞代碼。比如不要在 Executor 里用 requests 同步調(diào)用必須用 aiohttp不要直接調(diào) time.sleep必須用 await asyncio.sleep數(shù)據(jù)庫(kù)操作要么用異步驅(qū)動(dòng)要么把阻塞操作丟給線程池。初期我在這上面踩了不少坑比如某個(gè) Executor 里用了同步 SQLite 查詢結(jié)果整個(gè)事件循環(huán)被卡住所有消息處理都跟著變慢。7.2 單機(jī)壓測(cè)數(shù)據(jù)吞吐與延遲表現(xiàn)在一臺(tái) 4 核 8GB 的普通云服務(wù)器上我用 locust 對(duì) hermes-agent 做了壓測(cè)。模擬客戶端持續(xù) POST 請(qǐng)求每條消息大小約為 1KB。事件處理邏輯為空操作只做規(guī)則匹配和落庫(kù)。壓測(cè)結(jié)果穩(wěn)定狀態(tài)下每秒處理約 500 條消息P95 延遲為 60 毫秒P99 延遲為 120 毫秒。消息從接收到進(jìn)入 pending 隊(duì)列再被處理完成整體耗時(shí)平均在 300 毫秒左右。如果把規(guī)則匹配加上每條消息平均匹配 50 條規(guī)則吞吐量下降到 300 QPS但 P99 延遲仍在 200 毫秒以內(nèi)。對(duì)于內(nèi)部自動(dòng)化場(chǎng)景這個(gè)性能完全足夠。7.3 壓測(cè)中暴露的瓶頸與優(yōu)化過(guò)程壓測(cè)暴露的第一個(gè)瓶頸是 SQLite 單條插入性能。每條消息獨(dú)立 INSERT 提交在事務(wù)開(kāi)銷上浪費(fèi)了不少時(shí)間。優(yōu)化方案是批量插入攢一批消息后統(tǒng)一提交使寫(xiě)入吞吐量提升了一倍以上。這個(gè)批次大小我調(diào)到了 100 條一批平衡了延遲和吞吐。第二個(gè)瓶頸是規(guī)則深度匹配中的正則表達(dá)式操作。正則匹配雖然靈活但 CPU 開(kāi)銷明顯高于等值匹配。優(yōu)化方案是給規(guī)則增加條件預(yù)篩階段先用 JSONPath 取出的值做一次快速數(shù)據(jù)類型和長(zhǎng)度判斷過(guò)濾掉明顯不匹配的規(guī)則再進(jìn)入正則不歸。實(shí)際效果是把規(guī)則匹配整體耗時(shí)降低了 40%。第三個(gè)優(yōu)化是熱點(diǎn)路徑上的日志修改。壓測(cè)時(shí) statsd 和日志輸出占用了不少 CPU。我把每條消息處理的 DEBUG 日志改成按采樣率輸出比如每 100 條記錄一次整體性能提升明顯。對(duì)于日志這個(gè)點(diǎn)線上環(huán)境和壓測(cè)環(huán)境最大的區(qū)別就是線上大量無(wú)意義日志會(huì)把系統(tǒng)拖垮保留關(guān)鍵日志、關(guān)閉頻繁的 DEBUG 日志才是工程化做法。8. 部署與運(yùn)維實(shí)踐8.1 單機(jī) Docker 部署方案hermes-agent 的部署很簡(jiǎn)單。我提供了一個(gè) Dockerfile基于 python:3.11-slim 構(gòu)建體積控制在 200MB 以內(nèi)。整個(gè)鏡像只暴露一個(gè)端口通過(guò)環(huán)境變量注入配置。FROM python:3.11-slim WORKDIR /app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt COPY . . ENV HERMES_CONFIG/app/config/config.yaml EXPOSE 8080 CMD [python, -m, hermes_agent]啟動(dòng)命令也簡(jiǎn)單docker build -t hermes-agent:latest . docker run -d --name hermes-agent \ -p 8080:8080 \ -v /data/hermes:/app/data \ -e HERMES_CONFIG/app/config/config.yaml \ --restart unless-stopped \ hermes-agent:latest將數(shù)據(jù)目錄掛在宿主機(jī)上的好處是升級(jí)鏡像時(shí)數(shù)據(jù)不丟排錯(cuò)時(shí)可以直接查看 SQLite 文件和日志。8.2 配置管理環(huán)境變量與敏感信息配置管理這一塊我遵循的原則是非敏感配置放 YAML 文件敏感配置放環(huán)境變量。支付密鑰、數(shù)據(jù)庫(kù)密碼這類信息直接寫(xiě)在 YAML 里并且提交到代碼倉(cāng)庫(kù)是極其危險(xiǎn)的做法一旦代碼泄漏全部密鑰跟著泄漏。在實(shí)際部署中我通過(guò)環(huán)境變量注入這些敏感信息配置文件中只引用變量名。import os from dotenv import load_dotenv load_dotenv() PAYMENT_HOOK_SECRET os.getenv(PAYMENT_HOOK_SECRET)建議在 CI/CD 流水線中把生產(chǎn)環(huán)境的密鑰存儲(chǔ)在專門的密鑰管理服務(wù)中運(yùn)行時(shí)注入。即使是內(nèi)部系統(tǒng)的密鑰也不能因?yàn)橛X(jué)得沒(méi)那么重要就放松管理。8.3 日志規(guī)范與健康檢查接口hermes-agent 的日志統(tǒng)一輸出為 JSON 格式字段包括時(shí)間、級(jí)別、event_id、source、event_type、執(zhí)行器名稱、耗時(shí)、結(jié)果狀態(tài)。方便接入 ELK 或 Loki 做日志搜索。統(tǒng)一的 JSON 日志格式我堅(jiān)持了很久因?yàn)榕挪閱?wèn)題的時(shí)候跨字段關(guān)聯(lián)查詢比在純文本日志里用 grep 高效得多。健康檢查接口在/healthz上返回 200。檢查內(nèi)容包括進(jìn)程是否存活、SQLite 是否能正常讀寫(xiě)、掛起的 pending 消息數(shù)量是否超過(guò)閾值、最近 10 分鐘的錯(cuò)誤率等。K8s 或 Docker Compose 的健康檢查探針可以把它配置為探活地址。8.4 消息積壓與延遲的監(jiān)控指標(biāo)運(yùn)營(yíng) hermes-agent 的過(guò)程中我最關(guān)心的指標(biāo)有三個(gè)pending 消息數(shù)、處理延遲和執(zhí)行失敗率。這三個(gè)指標(biāo)分別對(duì)應(yīng)消息積壓、系統(tǒng)健康度和執(zhí)行質(zhì)量。pending 消息數(shù)可以通過(guò)一條 SQL 獲取SELECT status, COUNT(*) FROM events GROUP BY status;處理延遲指標(biāo)我通過(guò)記錄消息從 received_at 到 updated_at 的時(shí)間差在管理儀表盤上畫(huà)出分位數(shù)值。正常情況下 P95 應(yīng)該在 2 秒以下。如果某段時(shí)間 P95 突然上升基本可以判斷是某個(gè)執(zhí)行器變慢或者下游服務(wù)出現(xiàn)問(wèn)題。9. 項(xiàng)目演進(jìn)與多實(shí)例擴(kuò)展9.1 從 SQLite 遷移到 Postgres當(dāng)消息量增大到日均幾十萬(wàn)條時(shí)單機(jī) SQLite 開(kāi)始成為瓶頸。遷移到 Postgres 是自然的演進(jìn)方向。我在存儲(chǔ)層抽象了一個(gè)接口SQLite 和 Postgres 各實(shí)現(xiàn)一份切換時(shí)只需要改配置。表結(jié)構(gòu)幾乎不變只有數(shù)據(jù)類型上做了一些調(diào)整比如 TIMESTAMP 帶時(shí)區(qū)、JSON 字段用 JSONB 類型。storage: type: postgres dsn: postgresql://hermes:passwordlocalhost:5432/hermes9.2 多實(shí)例運(yùn)行時(shí)的任務(wù)鎖與消息分區(qū)Postgres 版本支持多實(shí)例部署也就是多個(gè) hermes-agent 進(jìn)程同時(shí)消費(fèi)同一個(gè)消息表。這里的關(guān)鍵問(wèn)題是不能讓兩個(gè)實(shí)例同時(shí)拿到同一條 pending 消息。解決辦法是使用 Postgres 的行級(jí)鎖通過(guò) SELECT FOR UPDATE SKIP LOCKED 實(shí)現(xiàn)安全的消息領(lǐng)取。SELECT * FROM events WHERE status pending ORDER BY priority DESC, created_at ASC LIMIT 1 FOR UPDATE SKIP LOCKED;這個(gè) SQL 的意思是鎖定并返回一條 pending 消息如果該行已經(jīng)被其他事務(wù)鎖定則跳過(guò)它不會(huì)阻塞等待。這讓多個(gè)實(shí)例可以并行安全地消費(fèi)消息無(wú)需額外引入分布式鎖。實(shí)際測(cè)試中3 個(gè)實(shí)例部署時(shí)吞吐量基本能線性擴(kuò)展到 1200 QPS。9.3 多執(zhí)行器協(xié)作的編排能力單條規(guī)則只能綁定一個(gè)動(dòng)作這個(gè)限制在復(fù)雜場(chǎng)景下會(huì)顯得不夠用。例如用戶下單后需要同時(shí)完成扣庫(kù)存、通知發(fā)貨、記錄財(cái)務(wù)流水這三個(gè)動(dòng)作如果失敗其中一個(gè)其他兩個(gè)應(yīng)該如何處理這就需要編排能力。為此我在新版本里引入了 action chain 的概念actions: - executor: inventory_executor params: operation: deduct - executor: financial_executor params: operation: record - executor: notify_executor params: channel: dingtalk template: order_created默認(rèn)情況下同一規(guī)則下的多個(gè)動(dòng)作按順序執(zhí)行前一個(gè)失敗則后續(xù)中斷。你也可以配置為并發(fā)執(zhí)行或者允許部分失敗繼續(xù)。這個(gè)編排能力讓 hermes-agent 從單純的消息轉(zhuǎn)發(fā)進(jìn)化為一個(gè)輕量級(jí)的任務(wù)調(diào)度引擎。9.4 多智能體協(xié)作模式的探索目前版本的 hermes-agent 支持讓執(zhí)行器反過(guò)來(lái)向消息隊(duì)列寫(xiě)入新事件這為多智能體協(xié)作提供了可能。例如一個(gè)訂單超時(shí)監(jiān)控執(zhí)行器發(fā)現(xiàn)訂單存在異常它可以生成一條新事件交給另一個(gè)專門處理異常的執(zhí)行器處理。Agent A 處理完自己的部分把結(jié)果作為新消息發(fā)布出來(lái)Agent B 接住再處理下一環(huán)。這樣的好處是處理鏈路上每個(gè)節(jié)點(diǎn)都清晰可追蹤天然支持分布式部署。這個(gè)模式我目前還在實(shí)踐中探索比如用于消息分片后的并行處理、跨部門業(yè)務(wù)鏈路的自動(dòng)化等場(chǎng)景。對(duì)于已經(jīng)上手的團(tuán)隊(duì)這可能是把 hermes-agent 能力翻倍的關(guān)鍵方向。10. 遇到問(wèn)題時(shí)的排查思路與實(shí)操經(jīng)驗(yàn)10.1 消息處理失敗但日志沒(méi)有明顯報(bào)錯(cuò)這種情況遇到過(guò)好幾次。排查的第一步是先確認(rèn)事件目前的狀態(tài)是 failed 還是 success如果是 failed去 events 表里看 retry_count 字段判斷是執(zhí)行器主動(dòng)返回失敗還是框架重試后放棄。第二步查看失敗執(zhí)行器的 message 字段框架會(huì)把 Executor 返回的失敗原因記錄進(jìn)去。第三步看 framework 日志重點(diǎn)查找該 event_id 對(duì)應(yīng)的執(zhí)行記錄和異常堆棧。這類問(wèn)題的根源往往是 Executor 內(nèi)部吞掉了異常返回了一個(gè)本該拋出的結(jié)果。實(shí)際排查一次之后定位到某個(gè)腳本里用了 try...except... 把異常吃掉并直接返回成功導(dǎo)致框架認(rèn)為執(zhí)行成功但實(shí)際上業(yè)務(wù)處理并沒(méi)有完成。這個(gè)教訓(xùn)后來(lái)讓我養(yǎng)成了一個(gè)習(xí)慣執(zhí)行器內(nèi)部絕不允許無(wú)差別捕獲異常后靜默通過(guò)必須顯式返回失敗結(jié)果。10.2 消息重復(fù)消費(fèi)問(wèn)題消息重復(fù)消費(fèi)的排查思路我每次都是先確認(rèn)事件表里 event_id 是否存在重復(fù)記錄。如果存在說(shuō)明消息源頭或接入層已經(jīng)產(chǎn)生了重復(fù)事件。如果不存在但業(yè)務(wù)層面處理了兩次說(shuō)明是執(zhí)行器冪等設(shè)計(jì)有問(wèn)題。一個(gè)具體的案例是某執(zhí)行器在調(diào)用下游接口時(shí)由于框架超時(shí)重試同一事件被執(zhí)行了兩次但兩次執(zhí)行都改了不同字段所以業(yè)務(wù)數(shù)據(jù)變得不一致。解決辦法是在執(zhí)行器里增加一個(gè)基于 event_id 的處理記錄表每次執(zhí)行前先檢查是否已經(jīng)執(zhí)行過(guò)。10.3 規(guī)則匹配結(jié)果和預(yù)期不符規(guī)則匹配問(wèn)題是日常排查里見(jiàn)得最多的一類。寫(xiě)規(guī)則的時(shí)候想的是應(yīng)該匹配這類消息結(jié)果實(shí)際運(yùn)行時(shí)不匹配或者錯(cuò)誤匹配了。排查步驟一般是先在管理界面查看實(shí)際消息的完整內(nèi)容確認(rèn) event_type、source 等關(guān)鍵字段值再打開(kāi) Debug 模式輸出規(guī)則匹配時(shí)的條件逐項(xiàng)判斷結(jié)果看具體是哪個(gè)條件不滿足或誤匹配。有一次排查到最后發(fā)現(xiàn)問(wèn)題竟然出在 YAML 配置文件里有人把鍵的縮進(jìn)寫(xiě)錯(cuò)了導(dǎo)致條件層級(jí)與預(yù)期不同匹配邏輯完全變了。這類問(wèn)題光看配置很難發(fā)現(xiàn)最好給規(guī)則文件加一個(gè)本地校驗(yàn)工具在 CI 階段檢查 YAML 格式和必需字段。11. 個(gè)人經(jīng)驗(yàn)總結(jié)與未來(lái)方向這個(gè)項(xiàng)目從最初的一個(gè)內(nèi)部腳本逐漸演化成了帶規(guī)則引擎、執(zhí)行器插件、重試機(jī)制和管理界面的輕量級(jí) agent 框架。整個(gè)過(guò)程中我最深的體會(huì)有三點(diǎn)。第一輕量是王道。能用一個(gè)進(jìn)程解決的事情就不要引入一整個(gè)微服務(wù)架構(gòu)。很多項(xiàng)目在初期規(guī)模很小的時(shí)候就把組件拆得滿天飛實(shí)際上每個(gè)組件都增加了排查問(wèn)題的復(fù)雜度。hermes-agent 選擇用 SQLite 起步、用 YAML 做配置這個(gè)能不依賴就不依賴的理念幫助它快速落地并且在真實(shí)場(chǎng)景中穩(wěn)定運(yùn)行。第二消息處理系統(tǒng)的核心價(jià)值在于可追蹤性。消息從哪兒來(lái)、被誰(shuí)處理、處理成功還是失敗、失敗原因是什么這些信息必須完整記錄。很多系統(tǒng)出問(wèn)題無(wú)法定位就是因?yàn)橄⑻幚礞溌泛诤谢?。我設(shè)計(jì) hermes-agent 時(shí)花費(fèi)了大量精力在事件狀態(tài)、執(zhí)行結(jié)果、日志格式這些看起來(lái)不性感但實(shí)際極其重要的細(xì)節(jié)上最終實(shí)踐證明這些都是值得的。第三冪等設(shè)計(jì)比重試機(jī)制更重要。沒(méi)有業(yè)務(wù)冪等做支撐重試機(jī)制反而會(huì)放大問(wèn)題。無(wú)論你用的是消息隊(duì)列、定時(shí)任務(wù)還是 Webhook都要養(yǎng)成同一事件可以重復(fù)執(zhí)行且結(jié)果一致的設(shè)計(jì)習(xí)慣。如果你也想自己實(shí)現(xiàn)一個(gè)類似的 agent 框架我建議你從最基礎(chǔ)的消息類型和狀態(tài)流轉(zhuǎn)開(kāi)始做先跑通一條完整鏈路再逐步加規(guī)則、執(zhí)行器、重試這些復(fù)雜度。直接在項(xiàng)目初期設(shè)計(jì)一個(gè)面條式的大框架容易讓自己陷入細(xì)節(jié)不可自拔。一個(gè)穩(wěn)定可靠的小系統(tǒng)遠(yuǎn)勝過(guò)一個(gè)看起來(lái)功能完整但處處漏風(fēng)的大架構(gòu)。