置 Redis 客戶端實(shí)現(xiàn)發(fā)布訂閱、限流與分布式鎖)
后端RPC框架API設(shè)計(jì)【免費(fèi)下載鏈接】orpcTypesafe APIs Made Simple 項(xiàng)目地址https://gitcode.com/gh_mirrors/or/orpc點(diǎn)擊查看免費(fèi)下載導(dǎo)讀orpc/bun是 oRPC 為 Bun 運(yùn)行時(shí)提供的適配器包它不引入任何第三方 Redis 依賴直接基于 Bun 自帶的RedisClient提供三類核心能力基于 Redis Pub/Sub 的發(fā)布訂閱BunRedisPublisher、基于固定窗口計(jì)數(shù)器的限流BunRedisRateLimiter以及基于SET NX PX的分布式鎖experimental_BunRedisLocker。讀完本文你將掌握這三個(gè)適配器的完整配置參數(shù)、底層實(shí)現(xiàn)原理、與 oRPC Procedure/Middleware 的集成方式以及如何通過倉庫內(nèi)的測試用例驗(yàn)證其在多進(jìn)程、多實(shí)例環(huán)境下的行為。orpc/bun 在 oRPC 生態(tài)中的定位oRPC 的口號是Typesafe APIs Made Simple整個(gè)項(xiàng)目圍繞類型安全的 API展開orpc/contract負(fù)責(zé)以契約作為單一事實(shí)來源orpc/server負(fù)責(zé)構(gòu)建 API 或?qū)崿F(xiàn)契約orpc/client負(fù)責(zé)端到端類型安全地消費(fèi) APIorpc/openapi為 API 增加 OpenAPI 兼容性見 packages/bun/README.md 的 Packages 表格。而orpc/bun屬于Framework ecosystem integrations分組官方定位是Adapters for Buns Redis —— 為 Publisher、Rate Limit、Lock 三個(gè)內(nèi)置能力提供基于 Bun Redis 的適配器。這與 packages/bun/package.json 中包的描述完全一致Bun integration for oRPC: Redis-backed pub/sub, rate limiting, and locking using Buns built-in Redis client也就是說orpc/bun自身不實(shí)現(xiàn)業(yè)務(wù)邏輯而是把 oRPC 三個(gè)通用能力發(fā)布訂閱、限流、鎖的存儲(chǔ)后端替換為 Bun 運(yùn)行時(shí)內(nèi)置的 Redis 客戶端從而讓 Bun 應(yīng)用在零額外依賴的前提下獲得跨進(jìn)程、跨實(shí)例共享狀態(tài)的能力。從 packages/bun/src/index.ts 可以看到包的公共 API 只有三個(gè)導(dǎo)出export * from ./redis-lock export * from ./redis-publisher export * from ./redis-ratelimit同時(shí)在 package.json 中它的運(yùn)行時(shí)依賴僅為orpc/client、orpc/experimental-lock、orpc/publisher、orpc/ratelimit、orpc/server、orpc/shared與standard-server/core——沒有任何 Redis 客戶端依賴因?yàn)?Redis 客戶端本身由 Bun 運(yùn)行時(shí)提供。前置條件與安裝使用orpc/bun需要滿足兩個(gè)前提運(yùn)行時(shí)必須是 Bun適配器直接操作bun模塊導(dǎo)出的RedisClient類型Bun 版本需包含其內(nèi)置 Redis 客戶端支持需要一個(gè)可連接的 Redis 服務(wù)測試腳本通過REDIS_URL環(huán)境變量指向 Redis 實(shí)例見下文測試與驗(yàn)證。安裝方式與 oRPC 其他 beta 包一致倉庫內(nèi)對應(yīng)命令為參考 apps/content/docs/helpers/publisher.mdx、apps/content/docs/helpers/ratelimit.mdx 與 apps/content/docs/helpers/lock.mdx 的 Installation 小節(jié)npm install orpc/bunbeta由于三個(gè)適配器依賴的通用能力分別位于orpc/publisher、orpc/ratelimit、orpc/experimental-lock中實(shí)際使用時(shí)一般還需要安裝對應(yīng)的能力包及其 schema 轉(zhuǎn)換包如orpc/zod。orpc/bun已將其聲明為自身依賴因此引入時(shí)會(huì)一并可用。BunRedisPublisher基于 Redis Pub/Sub 的跨進(jìn)程事件分發(fā)BunRedisPublisher是orpc/publisher的 Redis 適配器負(fù)責(zé)把事件通過 Redis Pub/Sub 分發(fā)到不同進(jìn)程的訂閱者并可選地借助 Redis Stream 實(shí)現(xiàn)錯(cuò)過的消息可補(bǔ)發(fā)resume?;居梅ê诵膶?shí)現(xiàn)位于 packages/bun/src/redis-publisher.ts它繼承自 packages/publisher/src/adapters/base-redis.ts 中的抽象基類BaseRedisPublisher泛型參數(shù)T extends Recordstring, object用于描述事件名到事件負(fù)載的類型映射import { BunRedisPublisher } from orpc/bun import { redis } from bun const publisher new BunRedisPublisher{ something-updated: { id: string } }(redis, { prefix: app:, resume: { enabled: true, seconds: 300, }, }) await publisher.publish(something-updated, { id: 123 }) // 回調(diào)式訂閱返回取消訂閱函數(shù) const unsubscribe await publisher.subscribe(something-updated, (payload) { console.log(payload.id) }) await unsubscribe()subscribe同時(shí)支持AsyncIterator 風(fēng)格可直接用于for await...of循環(huán)這在把事件流轉(zhuǎn)發(fā)給客戶端時(shí)非常有用詳見下文與 oRPC Procedure 集成。配置參數(shù)詳解BunRedisPublisher的完整選項(xiàng)由BunRedisPublisherOptionsredis-publisher.ts與BaseRedisPublisherOptionsbase-redis.ts共同定義參數(shù)默認(rèn)值說明subscriberredis.duplicate()首次訂閱時(shí)惰性創(chuàng)建專用于訂閱的 Redis 連接。因?yàn)镻ub/Sub 會(huì)接管連接處于訂閱狀態(tài)的客戶端無法再執(zhí)行普通命令因此必須使用獨(dú)立連接prefixRedis Key 與 Pub/Sub 頻道名的前綴用于多應(yīng)用共享同一 Redis 實(shí)例時(shí)的隔離serializernew RPCJsonSerializer()負(fù)載的序列化/反序列化器默認(rèn)使用orpc/client的 RPC JSON 序列化器支持自定義 handlers 序列化Date、自定義類等復(fù)雜類型resume.enabledfalse是否開啟事件補(bǔ)發(fā)。開啟后發(fā)布的事件會(huì)被臨時(shí)存入 Redis Stream新訂閱者可基于lastEventId從指定位置恢復(fù)resume.seconds3005 分鐘事件保留時(shí)長秒。出于性能考慮過期清理是惰性執(zhí)行的所以事件可能比該時(shí)長多存活一小段時(shí)間此外Publisher基類還提供maxBufferedEvents選項(xiàng)packages/publisher/src/publisher.ts控制 AsyncIterator 訂閱者的緩沖區(qū)上限默認(rèn)100設(shè)為0表示禁用緩沖、事件必須在下一條到達(dá)前被消費(fèi)設(shè)為1表示只保留最新事件適合實(shí)時(shí)狀態(tài)類場景設(shè)為Infinity則保留全部事件無丟失但內(nèi)存占用高。底層原理一條 Lua 腳本保證順序一致BaseRedisPublisher的發(fā)布邏輯base-redis.ts在開啟 resume 時(shí)不是先寫 Stream 再 PUBLISH而是通過一條原子 Lua 腳本PUBLISH_SCRIPT完成local idredis.call(XADD,KEYS[1],*,data,ARGV[1]) if ARGV[2] then redis.call(XTRIM,KEYS[1],MINID,ARGV[2],ARGV[3]) redis.call(EXPIRE,KEYS[1],ARGV[4]) end redis.call(PUBLISH,KEYS[1],{data:..ARGV[1]..,id:..id..})即XADD寫入 Stream → 按需XTRIM MINID裁剪過期條目并EXPIRE設(shè)置 TTL →PUBLISH到同名頻道。腳本的注釋明確解釋了設(shè)計(jì)動(dòng)機(jī)用一條腳本保證 Pub/Sub 投遞順序與 Stream 寫入順序一致從而避免補(bǔ)發(fā)與實(shí)時(shí)投遞之間出現(xiàn)順序錯(cuò)亂。BunRedisPublisher通過redis.send(EVAL, ...)redis-publisher.ts執(zhí)行該腳本并通過XREAD讀取補(bǔ)發(fā)數(shù)據(jù)redis-publisher.ts。在訂閱側(cè)base-redis.ts實(shí)現(xiàn)順序是先建立 Pub/Sub 訂閱再讀取 Stream 補(bǔ)發(fā)歷史最后處理訂閱期間積壓的實(shí)時(shí)消息并用resumedIds集合對補(bǔ)發(fā)與實(shí)時(shí)投遞之間可能競爭的重復(fù)事件做去重。這正是 redis-publisher.test.ts 中deduplicates events that race between resume and live delivery during reconnect用例所驗(yàn)證的行為。BunRedisRateLimiter固定窗口限流BunRedisRateLimiter是orpc/ratelimit的 Redis 適配器實(shí)現(xiàn)位于 packages/bun/src/redis-ratelimit.ts繼承自 packages/ratelimit/src/adapters/base-redis.ts 的BaseRedisRateLimiter。基本用法import { BunRedisRateLimiter } from orpc/bun import { redis } from bun const limiter new BunRedisRateLimiter(redis, { prefix: login:, maxRequests: 10, window: 60_000, // 毫秒 }) const result await limiter.limit(user:123, { weight: 2 }) if (!result.success) { // result 包含 limit / remaining / reset可據(jù)此拋出 ORPCError(TOO_MANY_REQUESTS, ...) }limit返回的完整結(jié)果為{ success, limit, remaining, reset }success表示本次請求是否被允許remaining為窗口內(nèi)剩余額度下限為 0reset是計(jì)數(shù)器重置的時(shí)間戳毫秒。配置參數(shù)參數(shù)默認(rèn)值說明prefixRedis Key 前綴用于隔離不同用途的計(jì)數(shù)器maxRequests無必填窗口內(nèi)允許的最大請求數(shù)window無必填固定窗口時(shí)長單位毫秒blockingUntilReady.enabledfalse是否開啟阻塞模式額度不足時(shí)等待而不是直接拒絕blockingUntilReady.timeout無enabled 時(shí)必填阻塞等待的最大時(shí)長毫秒超時(shí)后返回success: false底層原理原子 INCRBY 惰性窗口限流核心是一條固定窗口 Lua 腳本FIXED_WINDOW_SCRIPTlocal credis.call(INCRBY,KEYS[1],ARGV[1]) if ctonumber(ARGV[1]) then redis.call(PEXPIRE,KEYS[1],ARGV[2]) end return {c,redis.call(PTTL,KEYS[1])}即對計(jì)數(shù)器INCRBY增加本次請求的權(quán)重只有當(dāng)計(jì)數(shù)器是新建的c weight才設(shè)置PEXPIRE從而以惰性方式啟動(dòng)窗口最后返回[已用額度, 剩余 TTL]?;惖腸heckLimitbase-redis.ts據(jù)此計(jì)算reset Date.now() ttl。值得注意的兩個(gè)行為細(xì)節(jié)均有測試覆蓋見 redis-ratelimit.test.ts權(quán)重校驗(yàn)limit的weight必須是大于 0 的整數(shù)否則拋出TypeError(Rate limit weight must be an integer greater than 0)base-redis.ts阻塞模式blockUntilReadybase-redis.ts會(huì)循環(huán)調(diào)用checkLimit若失敗則sleep到reset時(shí)刻再試直到成功或超過timeout測試驗(yàn)證了它在窗口翻轉(zhuǎn)后放行加權(quán)請求、以及reset超出timeout時(shí)返回拒絕success: false兩種路徑。experimental_BunRedisLocker基于 SET NX PX 的分布式鎖experimental_BunRedisLocker是orpc/experimental-lock的 Redis 適配器類名帶experimental_前綴說明該能力仍處于實(shí)驗(yàn)階段。實(shí)現(xiàn)位于 packages/bun/src/redis-lock.ts繼承自 packages/lock/src/adapters/base-redis.ts 的BaseRedisLocker。基本用法import { experimental_BunRedisLocker as BunRedisLocker } from orpc/bun import { redis } from bun const locker new BunRedisLocker(redis, { ttl: 30_000, timeout: 5_000, }) const report await locker.lock(report:123, async ({ waited }) { // waited 為 true 表示曾等待其他持有者釋放鎖 return await generateReport(123) }, { ttl: 30_000, timeout: 5_000, signal: request.signal, })lock(key, fn, options)在持有鎖期間執(zhí)行回調(diào)回調(diào)拋出異常時(shí)也會(huì)釋放鎖finally保證見 base-redis.ts回調(diào)參數(shù)waited用于告知是否發(fā)生過等待——如果等待過說明同 key 的其他任務(wù)可能剛完成可先查緩存再重算。配置參數(shù)參數(shù)默認(rèn)值說明prefixRedis Key 前綴ttl無必填鎖的自動(dòng)過期時(shí)間毫秒防止持有者崩潰后鎖永不釋放可在每次調(diào)用時(shí)覆蓋timeout10000等待鎖可用的最長時(shí)間毫秒超時(shí)拋出LockTimeoutError可在調(diào)用時(shí)覆蓋retryInterval100鎖被他人持有時(shí)兩次獲取嘗試之間的間隔毫秒signal無可選 AbortSignal中止時(shí)提前結(jié)束等待并拋出中止原因僅在獲取鎖之前生效底層原理SET NX PX 加鎖 Lua 原子釋放加鎖使用 Redis 原生的SET key token NX PX ttlredis-lock.tsNX保證僅當(dāng) Key 不存在時(shí)才寫入即互斥PX設(shè)置過期時(shí)間返回OK表示獲取成功。鎖的 token 由crypto.randomUUID()生成base-redis.ts確保只有持有者本人能釋放鎖。釋放鎖不是簡單的DEL而是通過 Lua 腳本RELEASE_LOCK_SCRIPT先比對 token 再刪除避免持有者 A 的鎖已過期、被 B 重新獲取后A 卻把 B 的鎖刪掉的經(jīng)典問題if redis.call(GET, KEYS[1]) ARGV[1] then return redis.call(DEL, KEYS[1]) end return 0鎖的獲取循環(huán)base-redis.ts在未獲取成功時(shí)按retryInterval間隔重試直到timeout到期拋出LockTimeoutError(key)。與 oRPC Procedure / Middleware 的集成這三個(gè)適配器與 oRPC 的集成點(diǎn)在apps/content的 helpers 文檔中有完整示例。發(fā)布訂閱在 handler 中直接轉(zhuǎn)發(fā)事件流參考 apps/content/docs/helpers/publisher.mdx事件流可直接作為 Procedure 的輸出import { os } from orpc/server import * as z from zod const live os .handler(async function* ({ input, signal, lastEventId }) { const iterator publisher.subscribe(something-updated, { signal, lastEventId }) for await (const payload of iterator) { yield payload } }) const publish os .input(z.object({ id: z.string() })) .handler(async ({ input }) { await publisher.publish(something-updated, { id: input.id }) })倉庫的 playgrounds/bun/src 就是一個(gè)完整的 Bun 可運(yùn)行示例message.ts中subscribeMessages把publisher.subscribe(channel, { signal, lastEventId })直接作為asyncIteratorObject輸出返回客戶端即可實(shí)時(shí)收到消息。需要特別注意的是開啟 resume 后事件 id 由 publisher 自動(dòng)管理——發(fā)布時(shí)傳入的事件 id 會(huì)被忽略以 Redis Stream 分配的 id 為準(zhǔn)而服務(wù)端在 yield 自定義負(fù)載給客戶端時(shí)必須用getEventMeta(payload)?.id取出并隨withEventMeta透傳客戶端重連時(shí)才能正確傳回lastEventId續(xù)傳文檔以警告框形式強(qiáng)調(diào)了這一點(diǎn)。限流ratelimit 中間件參考 apps/content/docs/helpers/ratelimit.mdxratelimit中間件可基于 context 動(dòng)態(tài)選擇 limiterimport { ratelimit } from orpc/ratelimit const procedure os .$context{ ratelimiter: RateLimiter }() .input(z.object({ email: z.email() })) .use( ratelimit({ limiter: ({ context }) context.ratelimiter, key: ({ context }, input) login:${input.email}, weight: 1, // 每次請求消耗的額度默認(rèn) 1 }), ) .handler(({ input }) ({ success: true }))同一請求鏈中相同limiter key組合只會(huì)執(zhí)行一次限流檢查默認(rèn)去重配合RateLimitHandlerPlugin還能自動(dòng)在 HTTP 響應(yīng)中加入RateLimit-*與Retry-After響應(yīng)頭。鎖lock 中間件參考 apps/content/docs/helpers/lock.mdxlock中間件讓共享同一 key 的 Procedure 調(diào)用互斥執(zhí)行超時(shí)未獲取到鎖時(shí) Procedure 以CONFLICT錯(cuò)誤拒絕請求的signal會(huì)被轉(zhuǎn)發(fā)import { lock } from orpc/experimental-lock const procedure os .$context{ locker: Locker }() .input(z.object({ id: z.string() })) .use( lock({ locker: ({ context }) context.locker, key: ({ context }, input) report:${input.id}, ttl: 30_000, timeout: 5_000, }), ) .handler(async ({ context, input }) { if (context[lock/waited]) { // 同 key 的其他調(diào)用剛完成結(jié)果可能已可復(fù)用 } return await generateReport(input.id) })跨適配器兼容性與 Redis 客戶端適配器互通orpc/bun的三個(gè)適配器都不是獨(dú)立王國——它們與orpc/publisher、orpc/ratelimit、orpc/experimental-lock中基于 node-redis 的RedisPublisher/RedisRateLimiter/RedisLocker共享同一套 Redis Key 語義因此可以在同一集群中混用例如部分實(shí)例跑在 Bun 上、部分跑在 Node 上。這一點(diǎn)由三份兼容性測試用例背書packages/bun/tests/publisher-redis-adapters-compatibility.test.ts驗(yàn)證BunRedisPublisher與RedisPublisher互相投遞實(shí)時(shí)事件、并能從對方發(fā)布的 Stream 中按lastEventId續(xù)傳packages/bun/tests/ratelimit-redis-adapters-compatibility.test.ts驗(yàn)證二者共享限流計(jì)數(shù)器交替調(diào)用limit時(shí)剩余額度正確遞減、超限后success: falsepackages/bun/tests/lock-redis-adapters-compatibility.test.ts驗(yàn)證二者共享鎖狀態(tài)——一方持鎖時(shí)另一方等待釋放后等待者拿到waited: true。之所以能互通從源碼結(jié)構(gòu)看是因?yàn)槿齻€(gè) Bun 適配器都只實(shí)現(xiàn)了極薄的協(xié)議層真正復(fù)雜的 Key 命名規(guī)則、消息格式、Lua 腳本與重試/補(bǔ)發(fā)邏輯全部沉淀在各自的Base*抽象基類中適配器只負(fù)責(zé)把 BunRedisClient的命令調(diào)用映射為基類需要的四個(gè)原語publish、subscribe/unsubscribe、evalScript、readStreamEntries。測試與驗(yàn)證倉庫為orpc/bun提供了完整的集成測試運(yùn)行方式在 package.json 中定義bun --env-file../../.env test測試依賴真實(shí) Redis 服務(wù)通過REDIS_URL環(huán)境變量指定未設(shè)置時(shí)測試自動(dòng)跳過見各測試文件頂部的describe.skipIf(!REDIS_URL)注釋。測試分為三類單元/集成測試redis-publisher.test.ts、redis-ratelimit.test.ts、redis-lock.test.ts覆蓋事件補(bǔ)發(fā)順序、重連去重、并發(fā)發(fā)布者下的 Stream 順序、歷史裁剪與 TTL 過期、加權(quán)限流、阻塞模式、鎖超時(shí)LockTimeoutError、TTL 到期交接、AbortSignal 中止、以及同一 key 的回調(diào)絕不并發(fā)執(zhí)行測試斷言maxActive 1等關(guān)鍵行為跨適配器兼容測試上述tests/目錄下的三份*-redis-adapters-compatibility.test.tsRPC 傳輸測試packages/bun/tests/rpc/下還有基于 Bun 原生fetch與 WebSocket 的端到端傳輸測試如client-server.bun-fetch.ts、client-server.bun-websocket.ts用于驗(yàn)證 oRPC 服務(wù)在 Bun 運(yùn)行時(shí)下的完整鏈路。小結(jié)orpc/bun以極小的代碼面三個(gè)適配器把 Bun 內(nèi)置 Redis 客戶端接入 oRPC 的發(fā)布訂閱、限流與鎖體系底層通過統(tǒng)一的Base*基類與 Redis Lua 腳本保證原子性發(fā)布腳本保證 Pub/Sub 與 Stream 順序一致、限流腳本保證計(jì)數(shù)與窗口原子更新、釋放腳本保證只有 token 持有者可解鎖和跨適配器互通與 node-redis 適配器共享狀態(tài)。對于運(yùn)行在 Bun 上的 oRPC 應(yīng)用它是實(shí)現(xiàn)實(shí)時(shí)事件推送、接口限流、冪等任務(wù)互斥的首選零依賴方案。補(bǔ)充說明本文所述行為均以當(dāng)前倉庫代碼為準(zhǔn)orpc/bun版本2.0.0-beta.43API 可能隨 beta 迭代調(diào)整oRPC 官方文檔入口為 README 中引用的 packages/bun/README.md。oRPC 的設(shè)計(jì)靈感來自 tRPC端到端類型安全 RPC與 ts-rest契約優(yōu)先與 OpenAPI 集成這一點(diǎn)在 README 的 References 一節(jié)有明確致謝。贊分享后端RPC框架API設(shè)計(jì)【免費(fèi)下載鏈接】orpcTypesafe APIs Made Simple 項(xiàng)目地址https://gitcode.com/gh_mirrors/or/orpc點(diǎn)擊查看免費(fèi)下載上一篇自動(dòng)化治理架構(gòu)師實(shí)戰(zhàn)指南以 n8n 為中心的自動(dòng)化審計(jì)、風(fēng)險(xiǎn)評估與工作流治理下一篇如何用CSDN博客下載器實(shí)現(xiàn)技術(shù)知識體系化3步解決內(nèi)容碎片化難題創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考