workers: chain normalize jobs past the per-job cap and normalize periodically during long crawls
1 changed file +7 −1
modified
workers/main.ts
+7 −1
@@ -38,11 +38,15 @@ export async function startWorker(): Promise<() => Promise<void>> { | ||
| 38 | 38 | }); |
| 39 | 39 | await queue.work<{ connectorId?: string }>(JOBS.normalizeBatch, { concurrency: 2, pollingIntervalSeconds: 5 }, async (data) => { |
| 40 | 40 | let total = 0; |
| 41 | + let more = false; | |
| 41 | 42 | for (let i = 0; i < 40; i++) { |
| 42 | 43 | const r = await normalizeBatch({ connectorId: data.connectorId, limit: 500 }); |
| 43 | 44 | total += r.processed; |
| 44 | − if (r.processed < 500) break; | |
| 45 | + more = r.processed >= 500; | |
| 46 | + if (!more) break; | |
| 45 | 47 | } |
| 48 | + // Large crawls (100k+ raw rows) exceed one job's budget: chain another job instead of leaving rows unprocessed. | |
| 49 | + if (more) await queue.send(JOBS.normalizeBatch, { connectorId: data.connectorId }, { singletonKey: `normalize:${data.connectorId ?? 'all'}:${Date.now()}`, startAfterSeconds: 1 }); | |
| 46 | 50 | if (total > 0) for (let k = 0; k < Number(process.env.RESOLVE_CONCURRENCY ?? 3); k++) await queue.send(JOBS.resolveBatch, {}, { singletonKey: `resolve:${k}`, singletonSeconds: 30 }); |
| 47 | 51 | }); |
| 48 | 52 | const RESOLVE_WORKERS = Number(process.env.RESOLVE_CONCURRENCY ?? 3); |
@@ -112,6 +116,8 @@ export async function startWorker(): Promise<() => Promise<void>> { | ||
| 112 | 116 | const tick = async () => { |
| 113 | 117 | try { |
| 114 | 118 | await scheduleDueCrawls(queue); |
| 119 | + // Normalize whatever raw rows exist even while long crawls are still running (sales show up progressively). | |
| 120 | + await queue.send(JOBS.normalizeBatch, {}, { singletonKey: 'normalize:periodic', singletonSeconds: 120 }); | |
| 115 | 121 | } catch (err) { |
| 116 | 122 | log.error({ err }, 'scheduler tick failed'); |
| 117 | 123 | } |
| 118 | 124 | |