Skip to content

Add recurring jobs #4

Description

@yanatan16

Add recurring jobs to deetoo

The proposal is to add a recurring job type to deetoo. A recurring job is a delayed job that is requeued on the interval equal to the delay time.

Notating a Recurring Job

A delayed job is notated like so:

{
    "when": 1378828788,
    "id": "delayed_job_id" # optional
}

To notate a recurring job, the notation is as follows:

{
    "recurring": 300000, # recurring in milliseconds
    "id": "recurring_job_id", # required
    "now": true # optional, default true; Enqueues the job without a delay otherwise with delay equal to recurring
    "check_multiple": 3 # optional, default 3; Number of intervals to go before ensuring job is still alive.
}

Possible Issues

The main issue with a recurring job is ensuring the job is requeued and never lost. If we use the normal job execution system and append the requeuing at the end, there is nothing to ensure that requeue was successful. A Redis connection loss or partition can cause (and routinely does with Chinook) loss of requeue.

Proposed Solution

Upon "pushing" a recurring job in an application, deetoo creates a .when field as appropriate given .recurring and .now (and deleting .now in the processed, as its unnecessary). deetoo also starts a recurring function call that "ensures" the job is still there with an interval equal to .recurring * .check_multiple.

Note: Recurring jobs must be aligned with applications; therefore these jobs are suggested to be included during application startup.

Coded Pseudo-Solution

function pushRecurringJob(job, $next) {
    async.parallel([
        async.apply(createEnsureJob, job),
        function (cb) {
            if (!job.now) {
                requeueRecurringJob(job, cb);
            } else {
                Q.push(job, cb);
            }
        }
    ], $next);
}

function executeRecurringJob(job, $next) {
    async.parallel([
        async.apply(executeNormalJob, job),
        async.apply(requeueRecurringJob, job)
    ], $next);
}

function createEnsureJob(job, $next) {
    setInterval(ensure, job.recurring * (job.check_multiple || 3));

    function ensure () {
        kue.Job.get(job.id, function (err, kjob) {
            if (err) {
                return d2.log.error(err);
            }
            if (!job) {
                // Job is missing!? Let's add it back in
                d2.log.warn('Recurring job %s is missing.', job.id);
                requeueRecurringJob(job, function (err) {
                    if (err) {
                        d2.log.error(err);
                    }
                });
            }
        })
    }
}

function requeueRecurringJob(job, $next) {
    job.when = Date.now() + job.recurring;
    Q.push(job, $next);
}

Let me know any core flaws in my logic; and I'll code up a pull request soon that we can line-comment on.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions