feat(pipelines): 添加消息队列
This commit is contained in:
16
src/pipeline-tasks/pipeline-task.consumer.ts
Normal file
16
src/pipeline-tasks/pipeline-task.consumer.ts
Normal file
@@ -0,0 +1,16 @@
|
||||
import { OnQueueCompleted, Process, Processor } from '@nestjs/bull';
|
||||
import { Job } from 'bull';
|
||||
import { PipelineTask } from './pipeline-task.entity';
|
||||
import { PIPELINE_TASK_QUEUE } from './pipeline-tasks.constants';
|
||||
import { PipelineTasksService } from './pipeline-tasks.service';
|
||||
@Processor(PIPELINE_TASK_QUEUE)
|
||||
export class PipelineTaskConsumer {
|
||||
constructor(private readonly service: PipelineTasksService) {}
|
||||
@Process()
|
||||
async doTask() {}
|
||||
|
||||
@OnQueueCompleted()
|
||||
onCompleted(job: Job<PipelineTask>) {
|
||||
this.service.doNextTask(job.data.pipeline);
|
||||
}
|
||||
}
|
Reference in New Issue
Block a user