sql >> Base de Datos >  >> NoSQL >> Redis

Bull queue:cuando falla un trabajo, ¿cómo evitar que la cola procese los trabajos restantes?

En bull No es posible repetir el mismo trabajo inmediatamente después de su falla antes de continuar con el siguiente trabajo en la cola.

Solución:

  1. Cree un nuevo trabajo y establezca su prioridad en un valor menor que el tipo de trabajo actual.
  2. Libere el trabajo fallido (resolve() o done() )
  3. Este nuevo trabajo será recogido inmediatamente por el bull para su procesamiento.

Código de ejemplo:en el siguiente código Job-3 fallará y creará un nuevo trabajo y así sucesivamente hasta que el "propósito del trabajo" tenga éxito en algún momento.

var Queue = require('bull');

let redisOptions = {
  redis: { port: 6379, host: '127.0.0.1' }
}
var myQueue = new Queue('Linear-Queue', redisOptions);

myQueue.process('Type-1', function (job, done) {
  console.log(`Processing Job-${job.id} Attempt: ${job.attemptsMade}`);
  downloadFile(job, async function (error) {
    if (error) {
      await repeatSameJob(job, done);
    } else {
      done();
    }
  });
});

async function repeatSameJob(job, done) {
  let newJob = await myQueue.add('Type-1', job.data, { ...{ priority: 1 }, ...job.opts });
  console.log(`Job-${job.id} failed. Creating new Job-${newJob.id} with highest priority for same data.`);
  done(true);
}

function downloadFile(job, done) {
  setTimeout(async () => {
    done(job.data.error)
  }, job.data.time);
}

myQueue.on('completed', function (job, result) {
  console.log("Completed: Job-" + job.id);
});

myQueue.on('failed', async function (job, error) {
  console.log("Failed: Job-" + job.id);
});

let options = {
  removeOnComplete: true, // removes job from queue on success
  removeOnFail: true // removes job from queue on failure
}

for (let i = 1; i <= 5; i++) {
  let error = false;
  if (i == 3) { error = true; }

  setTimeout(i => {
    let jobData = {
      time: i * 2000,
      error: error,
      description: `Job-${i}`
    }
    myQueue.add('Type-1', jobData, options);
  }, i * 2000, i);
}

Salida: