Commit a91f8a37 by luoqi

fix(sync): 启动时回收僵尸同步锁,根治部署重启后 DW 滞后告警

现象:测试服 2026-07-23 8:00 报「DW 数据滞后」,增量游标 33.7h 未推进。

根因:sync_logs 用 partial UNIQUE(host_id) WHERE status='running' 作并发锁。
     部署 force-recreate 重启 service 容器时,正在跑的同步进程被杀,但那行
     running sync_log 来不及标终态 → 永久占住锁 → 之后每次 cron 增量都撞锁被
     skip → 游标不动。代码注释早写了"进程崩需依赖 stale 清理",但该清理从未实现。

修复:SyncIncrementalSchedulerService.onModuleInit 注册 cron 前,先把 startedAt
     早于本进程启动时刻(PROCESS_STARTED_AT)的 running 行标 failed。

判据安全性:sync 只在本 service 进程的 cron 回调、或一次性 CLI(跑完即退)里跑,
     两者都不可能比本进程启动得更早还活着。反过来,本进程启动后新建的 running
     (可能是并存的手动 CLI 真锁)用 lt 严格排除,绝不误清。
     幂等:标 failed 只释放锁,不动数据 —— 游标没推进,下次增量靠 48h 回看窗补齐。

数据零丢失:回看窗默认 48h > 本次空档 33.7h,靠 source_event_id 幂等去重完整补回。

测试:sync-lock-reap.spec 覆盖 回收/lt判据锁/空启动/失败不抛;295 全绿。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
parent 1f20bae4
Pipeline #3425 failed in 0 seconds
......@@ -2,11 +2,19 @@ import { Injectable, Logger, OnModuleInit } from '@nestjs/common';
import { SchedulerRegistry } from '@nestjs/schedule';
import { CronJob } from 'cron';
import * as path from 'path';
import { SyncStatus } from '@pac/types';
import { PrismaService } from '../prisma/prisma.service';
import {
ColdImportService,
SyncAlreadyRunningError,
} from '../modules/sync/cold-import/cold-import.service';
/**
* 本进程启动时刻(模块加载即固定)。
* 用作遗留锁判据:任何 startedAt 早于此刻的 running sync_log,其进程必然早于本进程启动 ——
* 而 sync 进程 = 本 service 进程(cron 内跑)或一次性 CLI(跑完即退)—— 都不可能还活着。
*/
const PROCESS_STARTED_AT = new Date();
import { OrgTreeService } from '../modules/auth/org-tree';
import { PersonaService } from '../modules/persona/persona.service';
import { PlanEngineService } from '../modules/plan/engine/plan-engine.service';
......@@ -41,7 +49,12 @@ export class SyncIncrementalSchedulerService implements OnModuleInit {
private readonly scheduler: SchedulerRegistry,
) {}
onModuleInit(): void {
async onModuleInit(): Promise<void> {
// 启动即回收上一个进程留下的僵尸锁(部署/崩溃重启时,running 行来不及标终态,
// 会永久占住 sync_logs 的 partial UNIQUE(host_id) WHERE status='running',
// 导致之后每次 cron 增量都被并发锁 skip → 游标不动 → DW 滞后告警)。
await this.reapStaleRunningLocks();
const globalCron = process.env.PAC_INCREMENTAL_CRON;
const dataDir = this.dataDir();
// env override(逗号分隔)优先;否则自动发现 auto_sync=true 的 host
......@@ -90,6 +103,46 @@ export class SyncIncrementalSchedulerService implements OnModuleInit {
return process.env.PAC_INCREMENTAL_DATA_DIR ?? path.resolve(__dirname, '../../data');
}
/**
* 回收僵尸同步锁 —— 把 startedAt 早于本进程启动的 running sync_log 标 failed。
*
* 为什么这样判据安全:sync 只在两处跑 —— 本 service 进程的 cron 回调,或一次性 CLI
* (`sync-incremental.cli` / `cold-import.cli`,跑完即退)。两者的进程都不可能比本进程
* 启动得更早还活着。所以 `startedAt < PROCESS_STARTED_AT` 的 running 行 = 上一个已死进程的残留。
* 反过来,本进程启动后新建的 running 行(startedAt >= PROCESS_STARTED_AT)绝不碰 —— 那可能是
* 正在跑的真锁(例如运维手动触发的 CLI 与本进程并存),误清会把在跑的同步" orphan"掉。
*
* 幂等:被标 failed 只是让并发锁释放;数据侧不受影响(游标没推进,下次增量靠 48h 回看窗补齐)。
*/
private async reapStaleRunningLocks(): Promise<void> {
try {
const stale = await this.prisma.syncLog.findMany({
where: { status: SyncStatus.RUNNING, startedAt: { lt: PROCESS_STARTED_AT } },
select: { id: true, hostId: true, startedAt: true, triggeredBy: true },
});
if (stale.length === 0) return;
const { count } = await this.prisma.syncLog.updateMany({
where: { status: SyncStatus.RUNNING, startedAt: { lt: PROCESS_STARTED_AT } },
data: {
status: SyncStatus.FAILED,
endedAt: new Date(),
errorMessage: '进程重启前未完成,启动时回收僵尸锁(reapStaleRunningLocks)',
},
});
for (const s of stale) {
this.logger.warn(
`sync-incremental: 回收僵尸锁 host=${s.hostId} started=${s.startedAt.toISOString()} ` +
`trigger=${s.triggeredBy}`,
);
}
this.logger.log(`sync-incremental: 启动回收 ${count} 个僵尸同步锁`);
} catch (err) {
// 回收失败不阻断启动(cron 仍注册);最坏退回人工清理,不比现状糟
this.logger.error(`sync-incremental: 僵尸锁回收失败: ${(err as Error).message}`);
}
}
/// 单 host 跑一轮(cron 回调用,吞异常不影响该 host 下次 / 别的 host)
private async runHostSafe(host: string): Promise<void> {
try {
......
import { SyncIncrementalSchedulerService } from '../src/queues/sync-incremental.scheduler';
/**
* 僵尸同步锁回收(reapStaleRunningLocks)回归。
*
* 背景:sync_logs 有 partial UNIQUE(host_id) WHERE status='running' 作并发锁。
* 部署/崩溃重启时,running 行来不及标终态,永久占住锁 → 之后每次 cron 增量被 skip →
* 游标不动 → DW 滞后告警(2026-07-23 测试服实测 33.7h)。
*
* 修复:service 启动时把 startedAt 早于本进程启动的 running 行标 failed。
* 判据安全性:sync 只在本进程 cron 或一次性 CLI 里跑,都不可能比本进程更早启动还活着;
* 本进程启动后新建的 running(可能是并存的手动 CLI 真锁)绝不碰。
*/
function makeService(rows: Array<{ id: string; hostId: string; startedAt: Date; triggeredBy: string }>) {
const findMany = jest.fn().mockImplementation(async ({ where }) => {
const lt: Date = where.startedAt.lt;
return rows.filter((r) => r.startedAt < lt);
});
const updateMany = jest.fn().mockImplementation(async ({ where }) => {
const lt: Date = where.startedAt.lt;
return { count: rows.filter((r) => r.startedAt < lt).length };
});
const prisma = { syncLog: { findMany, updateMany } };
const svc = new SyncIncrementalSchedulerService(
prisma as never,
{} as never,
{} as never,
{} as never,
{} as never,
{} as never,
);
return { svc, findMany, updateMany };
}
// 触达 private 方法(纯编排,无需走 onModuleInit 以免顺带注册 cron)
const reap = (svc: SyncIncrementalSchedulerService) =>
(svc as unknown as { reapStaleRunningLocks: () => Promise<void> }).reapStaleRunningLocks();
describe('reapStaleRunningLocks', () => {
test('回收进程启动前的 running 行 → 标 failed', async () => {
const old = new Date(Date.now() - 3_600_000); // 1h 前
const { svc, updateMany } = makeService([
{ id: 'zombie', hostId: 'h1', startedAt: old, triggeredBy: 'sync:jvs-dw:x' },
]);
await reap(svc);
expect(updateMany).toHaveBeenCalledTimes(1);
const arg = updateMany.mock.calls[0][0];
expect(arg.where.status).toBe('running');
expect(arg.data.status).toBe('failed');
expect(arg.data.endedAt).toBeInstanceOf(Date);
expect(typeof arg.data.errorMessage).toBe('string');
});
test('⭐ 只回收 startedAt < 进程启动时刻的行 —— where 用 lt,不用 gte', async () => {
// 判据锁:where 必须是 { startedAt: { lt: <某时刻> } },否则会误杀本进程启动后的真锁
const { svc, findMany, updateMany } = makeService([
{ id: 'zombie', hostId: 'h1', startedAt: new Date(Date.now() - 1000), triggeredBy: 't' },
]);
await reap(svc);
expect(findMany.mock.calls[0][0].where.startedAt).toHaveProperty('lt');
expect(findMany.mock.calls[0][0].where.startedAt).not.toHaveProperty('gte');
expect(updateMany.mock.calls[0][0].where.startedAt).toHaveProperty('lt');
});
test('无僵尸锁 → 不发 updateMany(纯净启动零副作用)', async () => {
const { svc, updateMany } = makeService([]);
await reap(svc);
expect(updateMany).not.toHaveBeenCalled();
});
test('回收失败不抛(不阻断启动)', async () => {
const prisma = {
syncLog: {
findMany: jest.fn().mockRejectedValue(new Error('db down')),
updateMany: jest.fn(),
},
};
const svc = new SyncIncrementalSchedulerService(
prisma as never, {} as never, {} as never, {} as never, {} as never, {} as never,
);
await expect(reap(svc)).resolves.toBeUndefined();
});
});
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment