ARTICLE DETAIL

资讯详情

深耕商务建站与企业官网运营的一线实战洞察。

oRPC 的 @orpc/bun 包:用 Bun 内置 Redis 客户端实现发布订阅、限流与分布式锁

oRPC 的 @orpc/bun 包:用 Bun 内置 Redis 客户端实现发布订阅、限流与分布式锁 后端RPC框架API设计【免费下载链接】orpcTypesafe APIs Made Simple 项目地址https://gitcode.com/gh_mirrors/or/orpc点击查看免费下载导读orpc/bun是 oRPC 为 Bun 运行时提供的适配器包它不引入任何第三方 Redis 依赖直接基于 Bun 自带的RedisClient提供三类核心能力基于 Redis Pub/Sub 的发布订阅BunRedisPublisher、基于固定窗口计数器的限流BunRedisRateLimiter以及基于SET NX PX的分布式锁experimental_BunRedisLocker。读完本文你将掌握这三个适配器的完整配置参数、底层实现原理、与 oRPC Procedure/Middleware 的集成方式以及如何通过仓库内的测试用例验证其在多进程、多实例环境下的行为。orpc/bun 在 oRPC 生态中的定位oRPC 的口号是Typesafe APIs Made Simple整个项目围绕类型安全的 API展开orpc/contract负责以契约作为单一事实来源orpc/server负责构建 API 或实现契约orpc/client负责端到端类型安全地消费 APIorpc/openapi为 API 增加 OpenAPI 兼容性见 packages/bun/README.md 的 Packages 表格。而orpc/bun属于Framework ecosystem integrations分组官方定位是Adapters for Buns Redis —— 为 Publisher、Rate Limit、Lock 三个内置能力提供基于 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自身不实现业务逻辑而是把 oRPC 三个通用能力发布订阅、限流、锁的存储后端替换为 Bun 运行时内置的 Redis 客户端从而让 Bun 应用在零额外依赖的前提下获得跨进程、跨实例共享状态的能力。从 packages/bun/src/index.ts 可以看到包的公共 API 只有三个导出export * from ./redis-lock export * from ./redis-publisher export * from ./redis-ratelimit同时在 package.json 中它的运行时依赖仅为orpc/client、orpc/experimental-lock、orpc/publisher、orpc/ratelimit、orpc/server、orpc/shared与standard-server/core——没有任何 Redis 客户端依赖因为 Redis 客户端本身由 Bun 运行时提供。前置条件与安装使用orpc/bun需要满足两个前提运行时必须是 Bun适配器直接操作bun模块导出的RedisClient类型Bun 版本需包含其内置 Redis 客户端支持需要一个可连接的 Redis 服务测试脚本通过REDIS_URL环境变量指向 Redis 实例见下文测试与验证。安装方式与 oRPC 其他 beta 包一致仓库内对应命令为参考 apps/content/docs/helpers/publisher.mdx、apps/content/docs/helpers/ratelimit.mdx 与 apps/content/docs/helpers/lock.mdx 的 Installation 小节npm install orpc/bunbeta由于三个适配器依赖的通用能力分别位于orpc/publisher、orpc/ratelimit、orpc/experimental-lock中实际使用时一般还需要安装对应的能力包及其 schema 转换包如orpc/zod。orpc/bun已将其声明为自身依赖因此引入时会一并可用。BunRedisPublisher基于 Redis Pub/Sub 的跨进程事件分发BunRedisPublisher是orpc/publisher的 Redis 适配器负责把事件通过 Redis Pub/Sub 分发到不同进程的订阅者并可选地借助 Redis Stream 实现错过的消息可补发resume。基本用法核心实现位于 packages/bun/src/redis-publisher.ts它继承自 packages/publisher/src/adapters/base-redis.ts 中的抽象基类BaseRedisPublisher泛型参数T extends Recordstring, object用于描述事件名到事件负载的类型映射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 }) // 回调式订阅返回取消订阅函数 const unsubscribe await publisher.subscribe(something-updated, (payload) { console.log(payload.id) }) await unsubscribe()subscribe同时支持AsyncIterator 风格可直接用于for await...of循环这在把事件流转发给客户端时非常有用详见下文与 oRPC Procedure 集成。配置参数详解BunRedisPublisher的完整选项由BunRedisPublisherOptionsredis-publisher.ts与BaseRedisPublisherOptionsbase-redis.ts共同定义参数默认值说明subscriberredis.duplicate()首次订阅时惰性创建专用于订阅的 Redis 连接。因为Pub/Sub 会接管连接处于订阅状态的客户端无法再执行普通命令因此必须使用独立连接prefixRedis Key 与 Pub/Sub 频道名的前缀用于多应用共享同一 Redis 实例时的隔离serializernew RPCJsonSerializer()负载的序列化/反序列化器默认使用orpc/client的 RPC JSON 序列化器支持自定义 handlers 序列化Date、自定义类等复杂类型resume.enabledfalse是否开启事件补发。开启后发布的事件会被临时存入 Redis Stream新订阅者可基于lastEventId从指定位置恢复resume.seconds3005 分钟事件保留时长秒。出于性能考虑过期清理是惰性执行的所以事件可能比该时长多存活一小段时间此外Publisher基类还提供maxBufferedEvents选项packages/publisher/src/publisher.ts控制 AsyncIterator 订阅者的缓冲区上限默认100设为0表示禁用缓冲、事件必须在下一条到达前被消费设为1表示只保留最新事件适合实时状态类场景设为Infinity则保留全部事件无丢失但内存占用高。底层原理一条 Lua 脚本保证顺序一致BaseRedisPublisher的发布逻辑base-redis.ts在开启 resume 时不是先写 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设置 TTL →PUBLISH到同名频道。脚本的注释明确解释了设计动机用一条脚本保证 Pub/Sub 投递顺序与 Stream 写入顺序一致从而避免补发与实时投递之间出现顺序错乱。BunRedisPublisher通过redis.send(EVAL, ...)redis-publisher.ts执行该脚本并通过XREAD读取补发数据redis-publisher.ts。在订阅侧base-redis.ts实现顺序是先建立 Pub/Sub 订阅再读取 Stream 补发历史最后处理订阅期间积压的实时消息并用resumedIds集合对补发与实时投递之间可能竞争的重复事件做去重。这正是 redis-publisher.test.ts 中deduplicates events that race between resume and live delivery during reconnect用例所验证的行为。BunRedisRateLimiter固定窗口限流BunRedisRateLimiter是orpc/ratelimit的 Redis 适配器实现位于 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可据此抛出 ORPCError(TOO_MANY_REQUESTS, ...) }limit返回的完整结果为{ success, limit, remaining, reset }success表示本次请求是否被允许remaining为窗口内剩余额度下限为 0reset是计数器重置的时间戳毫秒。配置参数参数默认值说明prefixRedis Key 前缀用于隔离不同用途的计数器maxRequests无必填窗口内允许的最大请求数window无必填固定窗口时长单位毫秒blockingUntilReady.enabledfalse是否开启阻塞模式额度不足时等待而不是直接拒绝blockingUntilReady.timeout无enabled 时必填阻塞等待的最大时长毫秒超时后返回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])}即对计数器INCRBY增加本次请求的权重只有当计数器是新建的c weight才设置PEXPIRE从而以惰性方式启动窗口最后返回[已用额度, 剩余 TTL]。基类的checkLimitbase-redis.ts据此计算reset Date.now() ttl。值得注意的两个行为细节均有测试覆盖见 redis-ratelimit.test.ts权重校验limit的weight必须是大于 0 的整数否则抛出TypeError(Rate limit weight must be an integer greater than 0)base-redis.ts阻塞模式blockUntilReadybase-redis.ts会循环调用checkLimit若失败则sleep到reset时刻再试直到成功或超过timeout测试验证了它在窗口翻转后放行加权请求、以及reset超出timeout时返回拒绝success: false两种路径。experimental_BunRedisLocker基于 SET NX PX 的分布式锁experimental_BunRedisLocker是orpc/experimental-lock的 Redis 适配器类名带experimental_前缀说明该能力仍处于实验阶段。实现位于 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)在持有锁期间执行回调回调抛出异常时也会释放锁finally保证见 base-redis.ts回调参数waited用于告知是否发生过等待——如果等待过说明同 key 的其他任务可能刚完成可先查缓存再重算。配置参数参数默认值说明prefixRedis Key 前缀ttl无必填锁的自动过期时间毫秒防止持有者崩溃后锁永不释放可在每次调用时覆盖timeout10000等待锁可用的最长时间毫秒超时抛出LockTimeoutError可在调用时覆盖retryInterval100锁被他人持有时两次获取尝试之间的间隔毫秒signal无可选 AbortSignal中止时提前结束等待并抛出中止原因仅在获取锁之前生效底层原理SET NX PX 加锁 Lua 原子释放加锁使用 Redis 原生的SET key token NX PX ttlredis-lock.tsNX保证仅当 Key 不存在时才写入即互斥PX设置过期时间返回OK表示获取成功。锁的 token 由crypto.randomUUID()生成base-redis.ts确保只有持有者本人能释放锁。释放锁不是简单的DEL而是通过 Lua 脚本RELEASE_LOCK_SCRIPT先比对 token 再删除避免持有者 A 的锁已过期、被 B 重新获取后A 却把 B 的锁删掉的经典问题if redis.call(GET, KEYS[1]) ARGV[1] then return redis.call(DEL, KEYS[1]) end return 0锁的获取循环base-redis.ts在未获取成功时按retryInterval间隔重试直到timeout到期抛出LockTimeoutError(key)。与 oRPC Procedure / Middleware 的集成这三个适配器与 oRPC 的集成点在apps/content的 helpers 文档中有完整示例。发布订阅在 handler 中直接转发事件流参考 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 就是一个完整的 Bun 可运行示例message.ts中subscribeMessages把publisher.subscribe(channel, { signal, lastEventId })直接作为asyncIteratorObject输出返回客户端即可实时收到消息。需要特别注意的是开启 resume 后事件 id 由 publisher 自动管理——发布时传入的事件 id 会被忽略以 Redis Stream 分配的 id 为准而服务端在 yield 自定义负载给客户端时必须用getEventMeta(payload)?.id取出并随withEventMeta透传客户端重连时才能正确传回lastEventId续传文档以警告框形式强调了这一点。限流ratelimit 中间件参考 apps/content/docs/helpers/ratelimit.mdxratelimit中间件可基于 context 动态选择 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, // 每次请求消耗的额度默认 1 }), ) .handler(({ input }) ({ success: true }))同一请求链中相同limiter key组合只会执行一次限流检查默认去重配合RateLimitHandlerPlugin还能自动在 HTTP 响应中加入RateLimit-*与Retry-After响应头。锁lock 中间件参考 apps/content/docs/helpers/lock.mdxlock中间件让共享同一 key 的 Procedure 调用互斥执行超时未获取到锁时 Procedure 以CONFLICT错误拒绝请求的signal会被转发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 的其他调用刚完成结果可能已可复用 } return await generateReport(input.id) })跨适配器兼容性与 Redis 客户端适配器互通orpc/bun的三个适配器都不是独立王国——它们与orpc/publisher、orpc/ratelimit、orpc/experimental-lock中基于 node-redis 的RedisPublisher/RedisRateLimiter/RedisLocker共享同一套 Redis Key 语义因此可以在同一集群中混用例如部分实例跑在 Bun 上、部分跑在 Node 上。这一点由三份兼容性测试用例背书packages/bun/tests/publisher-redis-adapters-compatibility.test.ts验证BunRedisPublisher与RedisPublisher互相投递实时事件、并能从对方发布的 Stream 中按lastEventId续传packages/bun/tests/ratelimit-redis-adapters-compatibility.test.ts验证二者共享限流计数器交替调用limit时剩余额度正确递减、超限后success: falsepackages/bun/tests/lock-redis-adapters-compatibility.test.ts验证二者共享锁状态——一方持锁时另一方等待释放后等待者拿到waited: true。之所以能互通从源码结构看是因为三个 Bun 适配器都只实现了极薄的协议层真正复杂的 Key 命名规则、消息格式、Lua 脚本与重试/补发逻辑全部沉淀在各自的Base*抽象基类中适配器只负责把 BunRedisClient的命令调用映射为基类需要的四个原语publish、subscribe/unsubscribe、evalScript、readStreamEntries。测试与验证仓库为orpc/bun提供了完整的集成测试运行方式在 package.json 中定义bun --env-file../../.env test测试依赖真实 Redis 服务通过REDIS_URL环境变量指定未设置时测试自动跳过见各测试文件顶部的describe.skipIf(!REDIS_URL)注释。测试分为三类单元/集成测试redis-publisher.test.ts、redis-ratelimit.test.ts、redis-lock.test.ts覆盖事件补发顺序、重连去重、并发发布者下的 Stream 顺序、历史裁剪与 TTL 过期、加权限流、阻塞模式、锁超时LockTimeoutError、TTL 到期交接、AbortSignal 中止、以及同一 key 的回调绝不并发执行测试断言maxActive 1等关键行为跨适配器兼容测试上述tests/目录下的三份*-redis-adapters-compatibility.test.tsRPC 传输测试packages/bun/tests/rpc/下还有基于 Bun 原生fetch与 WebSocket 的端到端传输测试如client-server.bun-fetch.ts、client-server.bun-websocket.ts用于验证 oRPC 服务在 Bun 运行时下的完整链路。小结orpc/bun以极小的代码面三个适配器把 Bun 内置 Redis 客户端接入 oRPC 的发布订阅、限流与锁体系底层通过统一的Base*基类与 Redis Lua 脚本保证原子性发布脚本保证 Pub/Sub 与 Stream 顺序一致、限流脚本保证计数与窗口原子更新、释放脚本保证只有 token 持有者可解锁和跨适配器互通与 node-redis 适配器共享状态。对于运行在 Bun 上的 oRPC 应用它是实现实时事件推送、接口限流、幂等任务互斥的首选零依赖方案。补充说明本文所述行为均以当前仓库代码为准orpc/bun版本2.0.0-beta.43API 可能随 beta 迭代调整oRPC 官方文档入口为 README 中引用的 packages/bun/README.md。oRPC 的设计灵感来自 tRPC端到端类型安全 RPC与 ts-rest契约优先与 OpenAPI 集成这一点在 README 的 References 一节有明确致谢。赞分享后端RPC框架API设计【免费下载链接】orpcTypesafe APIs Made Simple 项目地址https://gitcode.com/gh_mirrors/or/orpc点击查看免费下载上一篇自动化治理架构师实战指南以 n8n 为中心的自动化审计、风险评估与工作流治理下一篇如何用CSDN博客下载器实现技术知识体系化3步解决内容碎片化难题创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表
PREV
查看更多资讯
NEXT
返回资讯列表