perf: coalesce batch progress refreshes
This commit is contained in:
@@ -924,6 +924,28 @@ describe('SendChainService', () => {
|
|||||||
}));
|
}));
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('coalesces concurrent batch progress refreshes and keeps a trailing refresh', async () => {
|
||||||
|
const { service, prisma } = createService();
|
||||||
|
let resolveFirst: ((value: Array<{ status: string; _count: { _all: number } }>) => void) | undefined;
|
||||||
|
prisma.smsMessageRecord.groupBy
|
||||||
|
.mockImplementationOnce(() => new Promise((resolve) => { resolveFirst = resolve; }))
|
||||||
|
.mockResolvedValue([{ status: 'delivered', _count: { _all: 2 } }]);
|
||||||
|
|
||||||
|
const first = service['submission'].refreshTaskProgress('task-1');
|
||||||
|
await Promise.resolve();
|
||||||
|
const second = service['submission'].refreshTaskProgress('task-1');
|
||||||
|
|
||||||
|
expect(prisma.smsMessageRecord.groupBy).toHaveBeenCalledTimes(1);
|
||||||
|
resolveFirst?.([{ status: 'delivered', _count: { _all: 1 } }]);
|
||||||
|
await Promise.all([first, second]);
|
||||||
|
|
||||||
|
expect(prisma.smsMessageRecord.groupBy).toHaveBeenCalledTimes(2);
|
||||||
|
expect(prisma.smsBatchTask.update).toHaveBeenLastCalledWith({
|
||||||
|
where: { id: 'task-1' },
|
||||||
|
data: expect.objectContaining({ progressTotal: 2, successTotal: 2, status: 'finished' }),
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
it('does not expose CMPP internal tasks through client task detail or messages', async () => {
|
it('does not expose CMPP internal tasks through client task detail or messages', async () => {
|
||||||
const { service, prisma } = createService();
|
const { service, prisma } = createService();
|
||||||
prisma.smsBatchTask.findFirst.mockResolvedValue(null);
|
prisma.smsBatchTask.findFirst.mockResolvedValue(null);
|
||||||
|
|||||||
@@ -24,6 +24,8 @@ export class SendGatewaySubmitService {
|
|||||||
private sendQueue?: Queue<SendJob, unknown, 'send-message'>;
|
private sendQueue?: Queue<SendJob, unknown, 'send-message'>;
|
||||||
private gatewayQueue?: Queue;
|
private gatewayQueue?: Queue;
|
||||||
private worker?: Worker<SendJob>;
|
private worker?: Worker<SendJob>;
|
||||||
|
private readonly taskProgressRefreshes = new Map<string, Promise<void>>();
|
||||||
|
private readonly dirtyTaskProgressRefreshes = new Set<string>();
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
private readonly prisma: PrismaService,
|
private readonly prisma: PrismaService,
|
||||||
@@ -449,6 +451,29 @@ async waitForChannelRateLimit(channelId: string, tps: number) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async refreshTaskProgress(batchTaskId: string) {
|
async refreshTaskProgress(batchTaskId: string) {
|
||||||
|
const running = this.taskProgressRefreshes.get(batchTaskId);
|
||||||
|
if (running) {
|
||||||
|
// A state transition committed after the running aggregate may not be visible
|
||||||
|
// to its snapshot. Mark one trailing pass instead of issuing another identical
|
||||||
|
// GROUP BY concurrently for every message/result/receipt callback.
|
||||||
|
this.dirtyTaskProgressRefreshes.add(batchTaskId);
|
||||||
|
return running;
|
||||||
|
}
|
||||||
|
const refresh = this.refreshTaskProgressUntilClean(batchTaskId);
|
||||||
|
this.taskProgressRefreshes.set(batchTaskId, refresh);
|
||||||
|
try {
|
||||||
|
await refresh;
|
||||||
|
} finally {
|
||||||
|
if (this.taskProgressRefreshes.get(batchTaskId) === refresh) {
|
||||||
|
this.taskProgressRefreshes.delete(batchTaskId);
|
||||||
|
}
|
||||||
|
this.dirtyTaskProgressRefreshes.delete(batchTaskId);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private async refreshTaskProgressUntilClean(batchTaskId: string) {
|
||||||
|
do {
|
||||||
|
this.dirtyTaskProgressRefreshes.delete(batchTaskId);
|
||||||
const groups = await this.prisma.smsMessageRecord.groupBy({
|
const groups = await this.prisma.smsMessageRecord.groupBy({
|
||||||
by: ['status'],
|
by: ['status'],
|
||||||
where: { batchTaskId },
|
where: { batchTaskId },
|
||||||
@@ -468,6 +493,7 @@ async refreshTaskProgress(batchTaskId: string) {
|
|||||||
where: { id: batchTaskId },
|
where: { id: batchTaskId },
|
||||||
data: { progressTotal, submittedTotal, successTotal, failedTotal, unknownTotal, timeoutTotal, status },
|
data: { progressTotal, submittedTotal, successTotal, failedTotal, unknownTotal, timeoutTotal, status },
|
||||||
});
|
});
|
||||||
|
} while (this.dirtyTaskProgressRefreshes.has(batchTaskId));
|
||||||
}
|
}
|
||||||
|
|
||||||
getSendQueue(): Queue<SendJob, unknown, 'send-message'> {
|
getSendQueue(): Queue<SendJob, unknown, 'send-message'> {
|
||||||
|
|||||||
@@ -1172,3 +1172,4 @@ global,不能错误归入client。`client-signature-*`、发送页、企业认
|
|||||||
- `send-chain.helpers`提供无副作用的优先级、主备和稳定weight选择;提交服务只提供消息稳定键与真实在线/报备候选,不保存进程内轮询状态。
|
- `send-chain.helpers`提供无副作用的优先级、主备和稳定weight选择;提交服务只提供消息稳定键与真实在线/报备候选,不保存进程内轮询状态。
|
||||||
- Gateway既有Submit池与结果Outbox继续独立有界;先以正价六通道实测定位容量,只有证据证明单条回调仍触发停止线时才扩展批量回调协议。
|
- Gateway既有Submit池与结果Outbox继续独立有界;先以正价六通道实测定位容量,只有证据证明单条回调仍触发停止线时才扩展批量回调协议。
|
||||||
- Gateway上游模块对单分片使用聚合结果作为唯一Outbox边界,对多分片保留逐片持久化;结果Outbox不承担业务去重猜测,事件数量由上游已知分片数确定。
|
- Gateway上游模块对单分片使用聚合结果作为唯一Outbox边界,对多分片保留逐片持久化;结果Outbox不承担业务去重猜测,事件数量由上游已知分片数确定。
|
||||||
|
- 批次任务进度聚合仍归`send-gateway-submit`统一实现;同进程、同批次并发调用使用单飞与尾随刷新收敛重复`GROUP BY`,不改变消息状态写入、跨进程幂等或PostgreSQL最终事实。
|
||||||
|
|||||||
@@ -2136,4 +2136,5 @@
|
|||||||
- 供应商接受后,释放冻结与正式扣费以一个原子SQL写入两条幂等审计流水;二者净余额变化为零,不得争抢企业账户锁或重复更新账户。只有恢复历史“已释放但未扣费”半完成状态时才取得账户锁补扣。零金额不创建AccountTransaction;拒绝、超时、重投和最终失败继续沿用逐短信幂等释放、扣费和退款。
|
- 供应商接受后,释放冻结与正式扣费以一个原子SQL写入两条幂等审计流水;二者净余额变化为零,不得争抢企业账户锁或重复更新账户。只有恢复历史“已释放但未扣费”半完成状态时才取得账户锁补扣。零金额不创建AccountTransaction;拒绝、超时、重投和最终失败继续沿用逐短信幂等释放、扣费和退款。
|
||||||
- 同通道组内在线、已报备、同地域、同最低优先级且非备用的通道按weight和稳定消息键分流;备用、离线及更低优先级不抢占。只有显式配置为同优先级非备用的通道才参与主动双活。
|
- 同通道组内在线、已报备、同地域、同最低优先级且非备用的通道按weight和稳定消息键分流;备用、离线及更低优先级不抢占。只有显式配置为同优先级非备用的通道才参与主动双活。
|
||||||
- 单分片短短信的聚合SubmitResult已经携带完整分片信息,只发布聚合Outbox事件;不得再发布内容相同的分片事件造成双倍API回调和数据库写入。多分片长短信仍须在发送下一片前持久化当前分片事件,保持崩溃恢复边界。
|
- 单分片短短信的聚合SubmitResult已经携带完整分片信息,只发布聚合Outbox事件;不得再发布内容相同的分片事件造成双倍API回调和数据库写入。多分片长短信仍须在发送下一片前持久化当前分片事件,保持崩溃恢复边界。
|
||||||
|
- 同一批次并发提交结果、最终回执和失败处理触发任务进度刷新时,进程内只允许一个PostgreSQL聚合查询执行;并发触发合并为一次尾随刷新,既避免逐消息并发扫描整个批次,也必须覆盖运行中查询快照之后已经提交的状态变化。消息状态、计费和幂等事实仍以PostgreSQL为准,不得缓存聚合结果替代最终刷新。
|
||||||
- 压测使用隔离测试企业、真实PostgreSQL账务、单价`0.0325元/计费条`、三运营商混合号码和六个模拟供应商账号;逐档核对消息、冻结/释放/扣费/退款、SmsBillingRecord、通道分布、队列和数据库。临时单价与主动双活配置须先快照、测试后恢复。
|
- 压测使用隔离测试企业、真实PostgreSQL账务、单价`0.0325元/计费条`、三运营商混合号码和六个模拟供应商账号;逐档核对消息、冻结/释放/扣费/退款、SmsBillingRecord、通道分布、队列和数据库。临时单价与主动双活配置须先快照、测试后恢复。
|
||||||
|
|||||||
@@ -4786,3 +4786,4 @@ npm run verify:phase8
|
|||||||
| TC-CMPP-500-P4-007 | priority与normal并发积压 | priority保持明确服务能力且normal最终不饿死 |
|
| TC-CMPP-500-P4-007 | priority与normal并发积压 | priority保持明确服务能力且normal最终不饿死 |
|
||||||
| TC-CMPP-500-P4-008 | 测试配置治理 | 单价、账户与组项目先快照;临时主动双活和单价测试后恢复;预生产和凭据不变 |
|
| TC-CMPP-500-P4-008 | 测试配置治理 | 单价、账户与组项目先快照;临时主动双活和单价测试后恢复;预生产和凭据不变 |
|
||||||
| TC-CMPP-500-P4-009 | 单分片与多分片供应商Submit | 单分片只产生一个含segments的聚合Outbox事件;多分片逐片持久化并另有聚合事件,任何分片不丢失 |
|
| TC-CMPP-500-P4-009 | 单分片与多分片供应商Submit | 单分片只产生一个含segments的聚合Outbox事件;多分片逐片持久化并另有聚合事件,任何分片不丢失 |
|
||||||
|
| TC-CMPP-500-P4-010 | 同一批次并发任务进度刷新 | 并发调用共享一个运行中聚合查询,并在其间有新状态提交时只补一次尾随聚合;最终任务计数与消息终态一致 |
|
||||||
|
|||||||
@@ -3844,3 +3844,5 @@ git diff --check
|
|||||||
- 该轮最终双Stream归零,冻结/释放各2999、正式扣费2987、退款30、12条供应商最终拒绝只释放不扣费;SmsBillingRecord为charged2957/refunded30,余额恒等式精确成立,六通道与数据库最终无等待锁/idle in transaction。证据确认瓶颈是每个正价Submit结果在企业账户锁上串行,而不是单价0时可见的通道容量。
|
- 该轮最终双Stream归零,冻结/释放各2999、正式扣费2987、退款30、12条供应商最终拒绝只释放不扣费;SmsBillingRecord为charged2957/refunded30,余额恒等式精确成立,六通道与数据库最终无等待锁/idle in transaction。证据确认瓶颈是每个正价Submit结果在企业账户锁上串行,而不是单价0时可见的通道容量。
|
||||||
- 后续最小修复将正常“释放冻结+正式扣费”改为单个幂等SQL:冻结已在入口扣减,流水对净余额为零,因此正常结算不再取得账户锁或更新账户;历史已释放未扣费状态仍用原带锁路径补扣,并保留并发`ON CONFLICT`后一次MVCC恢复读。专项134项和TypeScript已通过,完整门禁、修复提交、新恢复资产与复测待补记。
|
- 后续最小修复将正常“释放冻结+正式扣费”改为单个幂等SQL:冻结已在入口扣减,流水对净余额为零,因此正常结算不再取得账户锁或更新账户;历史已释放未扣费状态仍用原带锁路径补扣,并保留并发`ON CONFLICT`后一次MVCC恢复读。专项134项和TypeScript已通过,完整门禁、修复提交、新恢复资产与复测待补记。
|
||||||
- 无账户锁修复后100条/秒复测仍在注入结束出现命令Stream`128/1687`、结果Outbox`32/452`;相同观测点完成扣费从388提高至637、等待锁降为0,但单分片短信同时发布分片与聚合两个结果事件,仍造成约双倍API回调和分片审计写入。新增最小Gateway修复:单分片只发布已携带完整segments的聚合事件,多分片仍逐片先持久化,待新恢复资产和同口径复测。
|
- 无账户锁修复后100条/秒复测仍在注入结束出现命令Stream`128/1687`、结果Outbox`32/452`;相同观测点完成扣费从388提高至637、等待锁降为0,但单分片短信同时发布分片与聚合两个结果事件,仍造成约双倍API回调和分片审计写入。新增最小Gateway修复:单分片只发布已携带完整segments的聚合事件,多分片仍逐片先持久化,待新恢复资产和同口径复测。
|
||||||
|
- 单分片回调收敛后第三轮100条/秒生成2999条,2999/2999受理,P50/P95/P99=`31/71/191ms`,三运营商`998/996/1005`、六通道`480至518`。注入期命令/结果Stream已排空但一度有2227条`submit_queued`,最终BullMQ wait/active/failed/delayed均0,消息2910 delivered、34 failed、55 submitted;冻结/释放各2999、charged 2987、refunded 21,SmsBillingRecord charged 2966/refunded 21,余额恒等式准确,数据库0等待锁/0 idle in transaction。该证据把剩余瓶颈定位到BullMQ发送Worker的逐消息数据库热路径,仍未达到全链100条/秒,因此没有升200档。
|
||||||
|
- 热路径确认每条提交结果和回执都会并发执行整批`SmsMessageRecord GROUP BY status`并更新任务。最小修复对同进程同批次刷新做单飞合并,并在运行中查询后出现新状态时保留一次尾随刷新,不缓存业务结果。新增`TC-CMPP-500-P4-010`;SendChain专项122项与API正式TypeScript编译通过,部署和同口径100档复测待补记。
|
||||||
|
|||||||
Reference in New Issue
Block a user