完整指南:单次请求发送 100 个事件与源码级原理解析)
Novu 批量触发Bulk Trigger完整指南单次请求发送 100 个事件与源码级原理解析【免费下载链接】novuThe open-source communication infrastructure for agents and products项目地址: https://gitcode.com/GitHub_Trending/no/novu本指南围绕 Novu 开源仓库中 bulk-trigger-examples.md 文档展开系统讲解如何通过triggerBulk在一次 API 调用中批量触发多个工作流事件、如何对大数量发送进行分块处理、如何使用 cURL 直连 REST 接口并结合仓库源码剖析批量触发在服务端的处理链路、100 事件上限的校验位置以及逐事件错误响应的实现机制。读完本文你将掌握用 TypeScript SDK 与 cURL 两种方式完成批量通知发送并能理解其底层实现从而在需要大规模、高效率触发通知时做出正确的技术选型。一、什么是 Bulk Trigger为什么需要批量触发在 Novu 中常规的POST /v1/events/trigger接口一次只能触发一个工作流事件虽然单个事件的to字段最多可携带 100 个收件人。当你的业务需要同时为大量订阅者触发不同的工作流例如向一批新用户发送欢迎邮件、向一批订单发送发货通知如果逐个调用单事件接口不仅会产生大量 HTTP 往返、增加延迟还容易触达 API 限流。Bulk Trigger批量触发正是为此设计的能力通过POST /v1/events/trigger/bulk接口可以在一次请求中携带最多 100 个独立事件每个事件可以对应不同的工作流、不同的订阅者、不同的 payload。从源码注释看该接口的定位非常明确——Using this endpoint you can trigger multiple events at once, to avoid multiple calls to the API见 events.controller.ts即用一次调用替代多次调用显著减少请求开销。从使用场景看批量触发特别适合以下场景新用户批量欢迎注册流程后一次性触达成百上千的新订阅者运营活动群发向不同用户群体发送内容各异的个性化消息系统事件回放把积压的业务事件一次性补偿触发多工作流组合发送同一批次里混合邮件、短信、App 内通知等不同类型的工作流。二、基本用法TypeScript SDK 一次发送多个事件文档给出的最简用法如下使用官方 TypeScript SDKnovu/apiimport { Novu } from novu/api; const novu new Novu({ secretKey: process.env.NOVU_SECRET_KEY, }); const result await novu.triggerBulk({ events: [ { workflowId: welcome-email, to: subscriber-1, payload: { userName: Alice }, }, { workflowId: welcome-email, to: subscriber-2, payload: { userName: Bob }, }, { workflowId: order-shipped, to: subscriber-3, payload: { orderId: ORD-100 }, }, ], });这段代码的关键点在于SDK 客户端初始化new Novu({ secretKey })使用环境变量NOVU_SECRET_KEY注入 API 密钥与单事件触发共用同一个客户端实例triggerBulk方法接收一个{ events: [...] }对象events是一个事件数组事件结构每个事件与单事件触发共享同一个请求体结构源码中为TriggerEventRequestDto包含workflowId、to、payload等核心字段事件彼此独立示例中前两个事件使用同一个welcome-email工作流发送给不同订阅者第三个事件则换成order-shipped工作流完全合法。值得说明的是SDK 中的triggerBulk方法名并非手工维护而是由服务端代码生成。在 events.controller.ts 中可以看到SdkMethodName(triggerBulk)与SdkUsageExample(Trigger Notification Events in Bulk)装饰器它们定义了生成 SDK 时的方法名与使用示例这也解释了为什么 TypeScript SDK 的方法签名与服务端 DTO 始终保持一致。每个事件的完整字段批量请求中的每个事件TriggerEventRequestDto定义于 trigger-event-request.dto.ts支持以下字段字段类型必填说明workflowId/namestring是工作流触发器标识符Trigger Identifier可在工作流页面找到SDK 中使用workflowIdREST 请求体中使用nametostring / string[] / 订阅者对象 / Topic 对象是收件人。支持订阅者 ID 字符串、订阅者对象含subscriberId、Topic 对象含topicKey与type单事件收件人数上限 100payloadobject否自定义数据用于渲染工作流模板内容、执行路由规则也会在通知 feed 中返回transactionIdstring否幂等去重标识相同transactionId再次触发会被忽略保留时长取决于计费套餐overridesobject否覆盖通道/步骤级配置例如指定 provider、替换 layoutsteps/channels/providers/severity等actorstring / 订阅者对象否显示在通知上的 Actor发送者头像tenantstring / Tenant 对象否指定租户上下文触发时可为已有租户更新信息contextobject否上下文信息最多 5 个bridgeUrlstring否可选的 Bridge Endpoint URL用于本地开发时把触发路由到指定 Bridge 应用必须是公网可访问的 https 地址agentIdstring / null否覆盖工作流默认绑定的 Agent传null可禁用本次执行的 Agent 路由三、100 事件上限与分块处理上限从哪里来DTO 与 Command 双重校验文档明确指出每次批量请求最多100 个事件。这个限制在服务端有两处落地请求体校验层在BulkTriggerEventDtotrigger-event-request.dto.ts中events数组使用ArrayNotEmpty()ValidateNested({ each: true })要求至少一个事件且逐项校验命令层校验ProcessBulkTriggerCommandprocess-bulk-trigger.command.ts中的ArrayMaxSize(100)明确限制数组最多 100 个元素。仓库的端到端测试 bulk-trigger.e2e.ts 专门验证了这一行为构造 101 个事件调用novuClient.triggerBulk断言返回 HTTP 422 且错误信息为events must contain no more than 100 elements。分块Chunking模式因此当事件数量超过 100 时必须在客户端将事件切分为多个批次依次发送。文档给出了一个通用的泛型分块函数function chunkT(array: T[], size: number): T[][] { const chunks: T[][] []; for (let i 0; i array.length; i size) { chunks.push(array.slice(i, i size)); } return chunks; } const allEvents users.map((user) ({ workflowId: weekly-digest, to: user.id, payload: { userName: user.name }, })); const batches chunk(allEvents, 100); for (const batch of batches) { await novu.triggerBulk({ events: batch }); }这段示例展示了两个非常实用的工程模式数据映射先从用户列表映射出事件数组users.map(...)每个事件对应一个收件人分块发送chunk(allEvents, 100)将大数组切成每块 100 个循环逐批调用triggerBulk保证每一批都符合服务端上限。如果要对数千甚至上万个收件人发送文档建议改用topic-based triggers基于 Topic 的触发先把订阅者添加进 Topic再对 Topic 触发一次事件服务端会向 Topic 内所有订阅者投递完全绕开单请求的事件数量限制适合超大规模群发场景。四、cURL 方式直连 REST 接口如果你的技术栈不使用官方 SDK也可以直接调用 REST 接口。文档给出的 curl 示例如下curl -X POST https://api.novu.co/v1/events/trigger/bulk \ -H Authorization: ApiKey $NOVU_SECRET_KEY \ -H Content-Type: application/json \ -d { events: [ { name: welcome-email, to: subscriber-1, payload: { userName: Alice } }, { name: welcome-email, to: subscriber-2, payload: { userName: Bob } } ] }与 SDK 调用对应的三个注意点请求方法POST路径为/v1/events/trigger/bulk对应服务端 events.controller.ts 中Post(/trigger/bulk)与全局/v1前缀的组合鉴权头Authorization: ApiKey $NOVU_SECRET_KEY密钥通过环境变量注入与 SDK 的secretKey配置等价字段名差异REST 请求体中工作流标识符字段叫name而 SDK 里叫workflowId——这是文档与源码中反复出现的两套命名本质指向同一个字段DTO 属性为nameSDK 生成时通过nameOverride: workflowId重命名见 trigger-event-request.dto.ts。响应体结构与 SDK 一致HTTP 201result数组按请求顺序返回每个事件的触发结果。五、源码级原理服务端如何批量处理这 100 个事件了解了客户端用法再看服务端实现。批量触发的核心处理逻辑位于 process-bulk-trigger.usecase.tsProcessBulkTrigger用例的执行过程可以拆成四步1. 批量预取工作流一次查询避免逐事件命中数据库const uniqueWorkflowIdentifiers [...new Set(command.events.map((event) event.name))]; const workflows await this.notificationTemplateRepository.find( { _environmentId: command.environmentId, triggers.identifier: { $in: uniqueWorkflowIdentifiers }, }, _id active payloadSchema validatePayload triggers, { readPreference: secondaryPreferred } );它先从所有事件中提取去重后的工作流标识符集合然后用 MongoDB 的$in一次性批量查询这些工作流只读取_id、active、payloadSchema、validatePayload、triggers等必要字段并指定secondaryPreferred读偏好以减轻主库压力最后构建workflowMap快速映射。这一步使得批量触发比 N 次单事件触发在数据库访问上高效得多。2. 内部再次分批每 5 个事件为一组并发解析拿到工作流后服务端并没有一次性并发处理全部 100 个事件而是使用BATCH_SIZE 5再次切片const BATCH_SIZE 5; for (let i 0; i command.events.length; i BATCH_SIZE) { const batch command.events.slice(i, i BATCH_SIZE); const batchResults await processBatch(batch); results.push(...batchResults); }每组内的 5 个事件通过Promise.all并发执行parseEventRequest事件解析与校验用例每组之间串行等待。这种外层每批 5 个的控制粒度可以看作对资源消耗的保守策略——避免 100 个事件瞬间并发造成 DB 与队列压力。3. 逐事件错误隔离一个失败不影响整批processBatch内每个事件都包在独立的try/catch中try { const workflow workflowMap.get(event.name); const result (await this.parseEventRequest.execute( ParseEventRequestMulticastCommand.create({ // ... 每个事件的字段被逐一透传 addressingType: AddressingTypeEnum.MULTICAST, requestCategory: TriggerRequestCategoryEnum.BULK, skipQueueInsertion: true, // ... }) )) as unknown as TriggerEventResponseDto; return result; } catch (e) { // 归一化错误信息返回带 error 的事件结果 return { acknowledged: true, status: TriggerEventStatusEnum.ERROR, error, transactionId: event.transactionId, } as TriggerEventResponseDto; }这正是文档中Errors/Success responses are returned per-event, not for the entire batch的源码出处某个事件失败例如workflow_not_found时只会让该事件在响应数组中带上status: error与error数组不会中断其他事件的正常处理。注意解析成功后设置了skipQueueInsertion: true说明事件先完成解析最后统一入队。4. 批量入队一次addBulk提交全部作业所有事件解析完毕后服务端把状态为processed且携带jobData的结果统一收集调用队列服务的批量入队接口const jobsToQueue: IWorkflowBulkJobDto[] results .filter((result) result.status TriggerEventStatusEnum.PROCESSED result.jobData ! undefined) .map((result) ({ name: result.jobData.transactionId, data: result.jobData, groupId: result.jobData.organizationId, })); if (jobsToQueue.length 0) { await this.workflowQueueService.addBulk(jobsToQueue); }最终返回时剥离jobData等内部字段把每个事件的公开结果按原顺序返回给调用方。整个设计呈现出一条清晰的流水线批量取模板 → 小批并发解析 → 逐事件错误隔离 → 批量入队。限流与鉴权批量端点还使用了独立的限流成本策略在 events.controller.ts 中标注了ThrottlerCost(ApiRateLimitCostEnum.BULK)而 apiRateLimits.ts 定义了SINGLE 1、BULK 100、KEYLESS 1000的成本权重——即一次批量触发按 100 次单事件触发的成本计入触发器类别的限流额度这与一次请求替代 100 次调用的定位完全对应。鉴权方面端点需要EVENT_WRITE权限并支持 External API Key、OAuth 与 Keyless 等访问方式。六、错误处理与响应语义响应结构按请求顺序返回的结果数组POST /v1/events/trigger/bulk返回 HTTP 201响应体为TriggerEventResponseDto[]见 events.controller.ts 的ApiResponse(TriggerEventResponseDto, 201, true)数组顺序与请求中events的顺序一致。每个事件的结果包含transactionId、statusprocessed或error、acknowledged等字段。仓库测试 bulk-trigger.e2e.ts 验证了这一点三个事件依次断言result.length 3且第 0/1/2 个结果的transactionId分别为1111/2222/3333状态均为processed、acknowledged为true。混合成功与失败逐事件错误bulk-trigger.e2e.ts 中有一个专门的用例三个事件中第一个指向不存在的non-existing-trigger第二个正常第三个缺少必填字段。结果是响应数组仍返回 3 个元素第一个事件status error且error[0] workflow_not_found第二个事件status processed。可见部分失败不会导致整批被拒绝你可以遍历响应数组、根据status与error字段对失败事件做重试或补偿。整体校验失败的场景如果请求本身不合法例如events为空数组、超过 100 个元素、某个事件 payload 不符合工作流 schema服务端会返回 4xx 错误——例如超过 100 元素返回 422payload 校验失败返回 400PayloadValidationExceptionDto提示 Payload validation failed - returned when any event payload does not match the workflow schema。这类错误属于整个请求层面的失败与逐事件的运行时错误不同需要区分对待。七、最佳实践总结综合文档要点与源码实现使用 Bulk Trigger 时的推荐做法如下单次最多 100 个事件客户端做好分块每块 100分块循环发送这是 DTO 与 Command 双重硬校验ArrayMaxSize(100)无法放宽利用事件独立性同一批次里混用不同工作流、不同订阅者、不同 payload 完全合法服务端会先批量去重加载工作流再逐事件解析与隔离错误逐事件消费响应不要假设整批成功或整批失败遍历result数组按status/error处理对error事件实施重试策略注意幂等transactionId可防止重复投递超大规模发送优先 Topic需要向数千以上订阅者群发时改用 Topic-based trigger——先建 Topic、添加订阅者再对 Topic 触发一次事件服务端统一投递比手动分块更省事也更高效留意限流成本一次批量触发按BULK 100的成本计入触发器限流类别批量并不能绕过限流而是把多次调用的开销合并到一次高频大批量场景仍要结合限流配额规划节奏。相关资源关联文档bulk-trigger-examples.md接口定义events.controller.ts请求体 DTOtrigger-event-request.dto.ts处理用例process-bulk-trigger.usecase.ts命令校验process-bulk-trigger.command.ts端到端测试bulk-trigger.e2e.ts限流成本定义apiRateLimits.ts【免费下载链接】novuThe open-source communication infrastructure for agents and products项目地址: https://gitcode.com/GitHub_Trending/no/novu创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考