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" }, }, ], });这段代码的关键点在于:
- 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/name | string | 是 | 工作流触发器标识符(Trigger Identifier),可在工作流页面找到;SDK 中使用workflowId,REST 请求体中使用name |
to | string / string[] / 订阅者对象 / Topic 对象 | 是 | 收件人。支持订阅者 ID 字符串、订阅者对象(含subscriberId)、Topic 对象(含topicKey与type),单事件收件人数上限 100 |
payload | object | 否 | 自定义数据,用于渲染工作流模板内容、执行路由规则,也会在通知 feed 中返回 |
transactionId | string | 否 | 幂等去重标识,相同transactionId再次触发会被忽略,保留时长取决于计费套餐 |
overrides | object | 否 | 覆盖通道/步骤级配置,例如指定 provider、替换 layout(steps/channels/providers/severity等) |
actor | string / 订阅者对象 | 否 | 显示在通知上的 Actor(发送者)头像 |
tenant | string / Tenant 对象 | 否 | 指定租户上下文,触发时可为已有租户更新信息 |
context | object | 否 | 上下文信息(最多 5 个) |
bridgeUrl | string | 否 | 可选的 Bridge Endpoint URL,用于本地开发时把触发路由到指定 Bridge 应用(必须是公网可访问的 https 地址) |
agentId | string / null | 否 | 覆盖工作流默认绑定的 Agent,传null可禁用本次执行的 Agent 路由 |
三、100 事件上限与分块处理
上限从哪里来:DTO 与 Command 双重校验
文档明确指出每次批量请求最多100 个事件。这个限制在服务端有两处落地:
- 请求体校验层:在
BulkTriggerEventDto(trigger-event-request.dto.ts)中,events数组使用@ArrayNotEmpty()+@ValidateNested({ each: true })要求至少一个事件且逐项校验; - 命令层校验:
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 调用对应的三个注意点:
- 请求方法:
POST,路径为/v1/events/trigger/bulk(对应服务端 events.controller.ts 中@Post('/trigger/bulk')与全局/v1前缀的组合); - 鉴权头:
Authorization: ApiKey $NOVU_SECRET_KEY,密钥通过环境变量注入,与 SDK 的secretKey配置等价; - 字段名差异: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一次性批量查询这些工作流(只读取_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、status(processed或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 元素返回 422,payload 校验失败返回 400(PayloadValidationExceptionDto,提示 "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
- 请求体 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),仅供参考