chore(server): adjust job config (#11405)

This commit is contained in:
forehalo
2025-04-02 09:48:45 +00:00
parent 85d176ce6f
commit 35bea20b80
5 changed files with 43 additions and 54 deletions

View File

@@ -56,15 +56,9 @@
"properties": { "properties": {
"apolloDriverConfig": { "apolloDriverConfig": {
"type": "object", "type": "object",
"description": "The config for underlying nestjs GraphQL and apollo driver engine.\n@default {\"buildSchemaOptions\":{\"numberScalarMode\":\"integer\"},\"useGlobalPrefix\":true,\"playground\":true,\"introspection\":true,\"sortSchema\":true}\n@link https://docs.nestjs.com/graphql/quick-start", "description": "The config for underlying nestjs GraphQL and apollo driver engine.\n@default {\"introspection\":true}\n@link https://docs.nestjs.com/graphql/quick-start",
"default": { "default": {
"buildSchemaOptions": { "introspection": true
"numberScalarMode": "integer"
},
"useGlobalPrefix": true,
"playground": true,
"introspection": true,
"sortSchema": true
} }
} }
} }
@@ -86,10 +80,8 @@
"properties": { "properties": {
"queue": { "queue": {
"type": "object", "type": "object",
"description": "The config for job queues\n@default {\"prefix\":\"affine_job\",\"defaultJobOptions\":{\"attempts\":5,\"removeOnComplete\":true,\"removeOnFail\":{\"age\":86400,\"count\":500}}}\n@link https://api.docs.bullmq.io/interfaces/v5.QueueOptions.html", "description": "The config for job queues\n@default {\"attempts\":5,\"removeOnComplete\":true,\"removeOnFail\":{\"age\":86400,\"count\":500}}\n@link https://api.docs.bullmq.io/interfaces/v5.QueueOptions.html",
"default": { "default": {
"prefix": "affine_job",
"defaultJobOptions": {
"attempts": 5, "attempts": 5,
"removeOnComplete": true, "removeOnComplete": true,
"removeOnFail": { "removeOnFail": {
@@ -97,7 +89,6 @@
"count": 500 "count": 500
} }
} }
}
}, },
"worker": { "worker": {
"type": "object", "type": "object",

View File

@@ -58,14 +58,12 @@ test.before(async () => {
ConfigModule.override({ ConfigModule.override({
job: { job: {
worker: { worker: {
defaultWorkerOptions: {
// NOTE(@forehalo): // NOTE(@forehalo):
// bullmq will hold the connection to check stalled jobs, // bullmq will hold the connection to check stalled jobs,
// which will keep the test process alive to timeout. // which will keep the test process alive to timeout.
stalledInterval: 100, stalledInterval: 100,
}, },
}, },
},
}), }),
JobModule.forRoot(), JobModule.forRoot(),
], ],

View File

@@ -1,4 +1,4 @@
import { QueueOptions, WorkerOptions } from 'bullmq'; import { DefaultJobOptions, WorkerOptions } from 'bullmq';
import { defineModuleConfig, JSONSchema } from '../../config'; import { defineModuleConfig, JSONSchema } from '../../config';
import { Queue } from './def'; import { Queue } from './def';
@@ -6,10 +6,8 @@ import { Queue } from './def';
declare global { declare global {
interface AppConfigSchema { interface AppConfigSchema {
job: { job: {
queue: ConfigItem<Omit<QueueOptions, 'connection' | 'telemetry'>>; queue: ConfigItem<Omit<DefaultJobOptions, 'connection' | 'telemetry'>>;
worker: ConfigItem<{ worker: ConfigItem<Omit<WorkerOptions, 'connection' | 'telemetry'>>;
defaultWorkerOptions: Omit<WorkerOptions, 'connection' | 'telemetry'>;
}>;
queues: { queues: {
[key in Queue]: ConfigItem<{ [key in Queue]: ConfigItem<{
concurrency: number; concurrency: number;
@@ -30,8 +28,6 @@ defineModuleConfig('job', {
queue: { queue: {
desc: 'The config for job queues', desc: 'The config for job queues',
default: { default: {
prefix: env.testing ? 'affine_job_test' : 'affine_job',
defaultJobOptions: {
attempts: 5, attempts: 5,
// should remove job after it's completed, because we will add a new job with the same job id // should remove job after it's completed, because we will add a new job with the same job id
removeOnComplete: true, removeOnComplete: true,
@@ -40,15 +36,12 @@ defineModuleConfig('job', {
count: 500, count: 500,
}, },
}, },
},
link: 'https://api.docs.bullmq.io/interfaces/v5.QueueOptions.html', link: 'https://api.docs.bullmq.io/interfaces/v5.QueueOptions.html',
}, },
worker: { worker: {
desc: 'The config for job workers', desc: 'The config for job workers',
default: { default: {},
defaultWorkerOptions: {},
},
link: 'https://api.docs.bullmq.io/interfaces/v5.WorkerOptions.html', link: 'https://api.docs.bullmq.io/interfaces/v5.WorkerOptions.html',
}, },

View File

@@ -1,7 +1,7 @@
import { getQueueToken } from '@nestjs/bullmq'; import { getQueueToken, getSharedConfigToken } from '@nestjs/bullmq';
import { Injectable, Logger, OnModuleDestroy } from '@nestjs/common'; import { Injectable, Logger, OnModuleDestroy } from '@nestjs/common';
import { ModuleRef } from '@nestjs/core'; import { ModuleRef } from '@nestjs/core';
import { Job, Queue as Bullmq, Worker } from 'bullmq'; import { Job, Queue as Bullmq, Worker, WorkerOptions } from 'bullmq';
import { difference, merge } from 'lodash-es'; import { difference, merge } from 'lodash-es';
import { CLS_ID, ClsServiceManager } from 'nestjs-cls'; import { CLS_ID, ClsServiceManager } from 'nestjs-cls';
@@ -118,16 +118,12 @@ export class JobExecutor implements OnModuleDestroy {
async job => { async job => {
return await this.run(job.name as JobName, job.data); return await this.run(job.name as JobName, job.data);
}, },
merge( merge({}, this.config.job.queue, this.config.job.worker, queueOptions, {
{}, prefix: this.ref.get(getSharedConfigToken(), { strict: false })
this.config.job.queue, .prefix,
this.config.job.worker.defaultWorkerOptions,
queueOptions,
{
concurrency, concurrency,
connection: this.redis, connection: this.redis,
} } as WorkerOptions)
)
); );
worker.on('error', error => { worker.on('error', error => {

View File

@@ -2,6 +2,7 @@ import './config';
import { BullModule } from '@nestjs/bullmq'; import { BullModule } from '@nestjs/bullmq';
import { DynamicModule } from '@nestjs/common'; import { DynamicModule } from '@nestjs/common';
import { type QueueOptions } from 'bullmq';
import { Config } from '../../config'; import { Config } from '../../config';
import { QueueRedis } from '../../redis'; import { QueueRedis } from '../../redis';
@@ -17,9 +18,19 @@ export class JobModule {
module: JobModule, module: JobModule,
imports: [ imports: [
BullModule.forRootAsync({ BullModule.forRootAsync({
useFactory: (config: Config, redis: QueueRedis) => { useFactory: (config: Config, redis: QueueRedis): QueueOptions => {
let prefix = 'affine_job';
if (env.testing) {
prefix += '_test';
} else if (!env.namespaces.production) {
prefix += '_' + env.NAMESPACE;
}
return { return {
...config.job.queue, // NOTE(@forehalo):
// we distinguish jobs by namespace,
// to avoid new jobs been dropped by old deployments
prefix,
defaultJobOptions: config.job.queue,
connection: redis, connection: redis,
}; };
}, },