news 2026/9/10 15:27:46

Novu 批量触发(Bulk Trigger)完整指南:单次请求发送 100 个事件与源码级原理解析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Novu 批量触发(Bulk Trigger)完整指南:单次请求发送 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 SDK@novu/api):

import { 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" }, }, ], });

这段代码的关键点在于:

  1. SDK 客户端初始化new Novu({ secretKey })使用环境变量NOVU_SECRET_KEY注入 API 密钥,与单事件触发共用同一个客户端实例;
  2. triggerBulk方法:接收一个{ events: [...] }对象,events是一个事件数组;
  3. 事件结构:每个事件与单事件触发共享同一个请求体结构(源码中为TriggerEventRequestDto),包含workflowIdtopayload等核心字段;
  4. 事件彼此独立:示例中前两个事件使用同一个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 中使用workflowId,REST 请求体中使用name
tostring / string[] / 订阅者对象 / Topic 对象收件人。支持订阅者 ID 字符串、订阅者对象(含subscriberId)、Topic 对象(含topicKeytype),单事件收件人数上限 100
payloadobject自定义数据,用于渲染工作流模板内容、执行路由规则,也会在通知 feed 中返回
transactionIdstring幂等去重标识,相同transactionId再次触发会被忽略,保留时长取决于计费套餐
overridesobject覆盖通道/步骤级配置,例如指定 provider、替换 layout(steps/channels/providers/severity等)
actorstring / 订阅者对象显示在通知上的 Actor(发送者)头像
tenantstring / Tenant 对象指定租户上下文,触发时可为已有租户更新信息
contextobject上下文信息(最多 5 个)
bridgeUrlstring可选的 Bridge Endpoint URL,用于本地开发时把触发路由到指定 Bridge 应用(必须是公网可访问的 https 地址)
agentIdstring / null覆盖工作流默认绑定的 Agent,传null可禁用本次执行的 Agent 路由

三、100 事件上限与分块处理

上限从哪里来:DTO 与 Command 双重校验

文档明确指出每次批量请求最多100 个事件。这个限制在服务端有两处落地:

  1. 请求体校验层:在BulkTriggerEventDto(trigger-event-request.dto.ts)中,events数组使用@ArrayNotEmpty()+@ValidateNested({ each: true })要求至少一个事件且逐项校验;
  2. 命令层校验ProcessBulkTriggerCommand(process-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 chunk<T>(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 调用对应的三个注意点:

  1. 请求方法POST,路径为/v1/events/trigger/bulk(对应服务端 events.controller.ts 中@Post('/trigger/bulk')与全局/v1前缀的组合);
  2. 鉴权头Authorization: ApiKey $NOVU_SECRET_KEY,密钥通过环境变量注入,与 SDK 的secretKey配置等价;
  3. 字段名差异:REST 请求体中工作流标识符字段叫name,而 SDK 里叫workflowId——这是文档与源码中反复出现的两套命名,本质指向同一个字段(DTO 属性为name,SDK 生成时通过nameOverride: 'workflowId'重命名,见 trigger-event-request.dto.ts)。

响应体结构与 SDK 一致:HTTP 201,result数组按请求顺序返回每个事件的触发结果。

五、源码级原理:服务端如何批量处理这 100 个事件

了解了客户端用法,再看服务端实现。批量触发的核心处理逻辑位于 process-bulk-trigger.usecase.ts,ProcessBulkTrigger用例的执行过程可以拆成四步:

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一次性批量查询这些工作流(只读取_idactivepayloadSchemavalidatePayloadtriggers等必要字段,并指定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 = 1BULK = 100KEYLESS = 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的顺序一致。每个事件的结果包含transactionIdstatusprocessederror)、acknowledged等字段。

仓库测试 bulk-trigger.e2e.ts 验证了这一点:三个事件依次断言result.length === 3,且第 0/1/2 个结果的transactionId分别为1111/2222/3333,状态均为processedacknowledgedtrue

混合成功与失败:逐事件错误

bulk-trigger.e2e.ts 中有一个专门的用例:三个事件中第一个指向不存在的non-existing-trigger,第二个正常,第三个缺少必填字段。结果是:响应数组仍返回 3 个元素,第一个事件status === 'error'error[0] === 'workflow_not_found',第二个事件status === 'processed'。可见部分失败不会导致整批被拒绝,你可以遍历响应数组、根据statuserror字段对失败事件做重试或补偿。

整体校验失败的场景

如果请求本身不合法(例如events为空数组、超过 100 个元素、某个事件 payload 不符合工作流 schema),服务端会返回 4xx 错误——例如超过 100 元素返回 422,payload 校验失败返回 400(PayloadValidationExceptionDto,提示 "Payload validation failed - returned when any event payload does not match the workflow schema")。这类错误属于整个请求层面的失败,与逐事件的运行时错误不同,需要区分对待。

七、最佳实践总结

综合文档要点与源码实现,使用 Bulk Trigger 时的推荐做法如下:

  1. 单次最多 100 个事件:客户端做好分块(每块 100),分块循环发送;这是 DTO 与 Command 双重硬校验(@ArrayMaxSize(100)),无法放宽;
  2. 利用事件独立性:同一批次里混用不同工作流、不同订阅者、不同 payload 完全合法;服务端会先批量去重加载工作流,再逐事件解析与隔离错误;
  3. 逐事件消费响应:不要假设整批成功或整批失败,遍历result数组按status/error处理,对error事件实施重试策略(注意幂等transactionId可防止重复投递);
  4. 超大规模发送优先 Topic:需要向数千以上订阅者群发时,改用 Topic-based trigger——先建 Topic、添加订阅者,再对 Topic 触发一次事件,服务端统一投递,比手动分块更省事也更高效;
  5. 留意限流成本:一次批量触发按BULK = 100的成本计入触发器限流类别,批量并不能"绕过"限流,而是把多次调用的开销合并到一次;高频大批量场景仍要结合限流配额规划节奏。

相关资源

  • 关联文档:bulk-trigger-examples.md
  • 接口定义:events.controller.ts
  • 请求体 DTO:trigger-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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/10 15:18:10

Buzz 离线语音转写怎么快速配通?Faster-Whisper 三步跑起来

Buzz 离线语音转写怎么快速配通&#xff1f;Faster-Whisper 三步跑起来 【免费下载链接】buzz Buzz transcribes and translates audio offline on your personal computer. Powered by OpenAIs Whisper. 项目地址: https://gitcode.com/GitHub_Trending/buz/buzz Buzz …

作者头像 李华