| 1 | import type { DynamicModule } from '@nestjs/common' |
| 2 | import { BullModule } from '@nestjs/bullmq' |
| 3 | import { Module } from '@nestjs/common' |
| 4 | import { RedisModule } from '@yikart/redis' |
| 5 | import { Redis } from 'ioredis' |
| 6 | import { QueueName } from './enums' |
| 7 | import { QueueMetricsService } from './queue-metrics.service' |
| 8 | import { QueueConfig } from './queue.config' |
| 9 | import { QueueService } from './queue.service' |
| 10 | import { createPinoTelemetry } from './telemetry/pino-telemetry' |
| 11 | |
| 12 | /** |
| 13 | * 队列模块 |
| 14 | * 提供统一的队列管理功能 |
| 15 | */ |
| 16 | @Module({}) |
| 17 | export class AitoearnQueueModule { |
| 18 | /** |
| 19 | * 创建队列模块 |
| 20 | * @param config 队列配置 |
| 21 | */ |
| 22 | static forRoot(config: QueueConfig): DynamicModule { |
| 23 | // 动态注册所有队列 |
| 24 | const queueModules = Object.values(QueueName).map((name) => { |
| 25 | // 可以在这里为特殊队列添加特殊配置 |
| 26 | // if (name === QueueName.SomeSpecialQueue) { |
| 27 | // return BullModule.registerQueue({ |
| 28 | // name, |
| 29 | // limiter: { duration: 10000, max: 100 }, |
| 30 | // }) |
| 31 | // } |
| 32 | return BullModule.registerQueue({ name }) |
| 33 | }) |
| 34 | |
| 35 | return { |
| 36 | global: true, |
| 37 | module: AitoearnQueueModule, |
| 38 | imports: [ |
| 39 | // 集成 Redis 模块 |
| 40 | RedisModule.forRoot(config.redis), |
| 41 | // 配置 Bull 连接 |
| 42 | BullModule.forRootAsync({ |
| 43 | imports: [], |
| 44 | useFactory: (redis: Redis) => ({ |
| 45 | prefix: config.prefix, |
| 46 | connection: redis, |
| 47 | telemetry: createPinoTelemetry(), |
| 48 | }), |
| 49 | inject: [Redis], |
| 50 | }), |
| 51 | // 注册所有队列 |
| 52 | ...queueModules, |
| 53 | ], |
| 54 | providers: [ |
| 55 | { |
| 56 | provide: QueueConfig, |
| 57 | useValue: config, |
| 58 | }, |
| 59 | QueueService, |
| 60 | QueueMetricsService, |
| 61 | ], |
| 62 | exports: [QueueService], |
| 63 | } |
| 64 | } |
| 65 | } |
| 66 |