Sign inSign up

graphile/worker

By graphile

•Updated 3 months ago

High performance Node.js/PostgreSQL job queue

Image
1

10K+

graphile/worker repository overview

⁠graphile-worker

Patreon sponsor button Discord chat room Package on npm MIT license Follow

Job queue for PostgreSQL running on Node.js - allows you to run jobs (e.g. sending emails, performing calculations, generating PDFs, etc) "in the background" so that your HTTP response/application code is not held up. Can be used with any PostgreSQL-backed application. Pairs beautifully with PostGraphile⁠ or PostgREST⁠.

⁠Crowd-funded open-source software

To help us develop this software sustainably under the MIT license, we ask all individuals and businesses that use it to help support its ongoing maintenance and development via sponsorship.

⁠Click here to find out more about sponsors and sponsorship.⁠

And please give some love to our featured sponsors 🤩:

* Sponsors the entire Graphile suite

⁠Quickstart: CLI

In your existing Node.js project:

⁠Add the worker to your project:
yarn add graphile-worker
# or: npm install --save graphile-worker
⁠Create tasks:

Create a tasks/ folder, and place in it JS files containing your task specs. The names of these files will be the task identifiers, e.g. hello below:

// tasks/hello.js
module.exports = async (payload, helpers) => {
  const { name } = payload;
  helpers.logger.info(`Hello, ${name}`);
};
⁠Run the worker

(Make sure you're in the folder that contains the tasks/ folder.)

npx graphile-worker -c "my_db"
# or, if you have a remote database, something like:
#   npx graphile-worker -c "postgres://user:pass@host:port/db?ssl=true"
# or, if you prefer envvars
#   DATABASE_URL="..." npx graphile-worker

(Note: npx runs the local copy of an npm module if it is installed, when you're ready, switch to using the package.json "scripts" entry instead.)

⁠Schedule a job via SQL

Connect to your database and run the following SQL:

SELECT graphile_worker.add_job('hello', json_build_object('name', 'Bobby Tables'));
⁠Success!

You should see the worker output Hello, Bobby Tables. Gosh, that was fast!

⁠Quickstart: library

Instead of running graphile-worker via the CLI, you may use it directly in your Node.js code. The following is equivalent to the CLI example above:

const { run, quickAddJob } = require("graphile-worker");

async function main() {
  // Run a worker to execute jobs:
  const runner = await run({
    connectionString: "postgres:///my_db",
    concurrency: 5,
    // Install signal handlers for graceful shutdown on SIGINT, SIGTERM, etc
    noHandleSignals: false,
    pollInterval: 1000,
    // you can set the taskList or taskDirectory but not both
    taskList: {
      hello: async (payload, helpers) => {
        const { name } = payload;
        helpers.logger.info(`Hello, ${name}`);
      },
    },
    // or:
    //   taskDirectory: `${__dirname}/tasks`,
  });

  // Or add a job to be executed:
  await quickAddJob(
    // makeWorkerUtils options
    { connectionString: "postgres:///my_db" },

    // Task identifier
    "hello",

    // Payload
    { name: "Bobby Tables" },
  );

  // If the worker exits (whether through fatal error or otherwise), this
  // promise will resolve/reject:
  await runner.promise;
}

main().catch((err) => {
  console.error(err);
  process.exit(1);
});

Running this example should output something like:

[core] INFO: Worker connected and looking for jobs... (task names: 'hello')
[job(worker-7327280603017288: hello{1})] INFO: Hello, Bobby Tables
[worker(worker-7327280603017288)] INFO: Completed task 1 (hello) with success (0.16ms)

⁠Support

You can ask for help on Discord at http://discord.gg/graphile⁠

Please support development of this project via sponsorship⁠. With your support we can improve performance, usability and documentation at a greater rate, leading to reduced running and engineering costs for your organisation, leading to a net ROI.

Professional support contracts are also available; for more information see: https://graphile.org/support/⁠

⁠Features

  • Standalone and embedded modes
  • Designed to be used both from JavaScript or directly in the database
  • Easy to test (recommended: runTaskListOnce util)
  • Low latency (typically under 3ms from task schedule to execution, uses LISTEN/NOTIFY to be informed of jobs as they're inserted)
  • High performance (uses SKIP LOCKED to find jobs to execute, resulting in faster fetches)
  • Small tasks (uses explicit task names / payloads resulting in minimal serialisation/deserialisation overhead)
  • Parallel by default
  • Adding jobs to same named queue runs them in series
  • Automatically re-attempts failed jobs with exponential back-off
  • Customisable retry count (default: 25 attempts over ~3 days)
  • Crontab-like scheduling feature for recurring tasks (with optional backfill)
  • Task de-duplication via unique job_key
  • Flexible runtime controls that can be used for complex rate limiting (e.g. via (graphile-worker-rate-limiter)[https://github.com/politics-rewired/graphile-worker-rate-limiter⁠])
  • Open source; liberal MIT license
  • Executes tasks written in Node.js (these can call out to any other language or networked service)
  • Modern JS with 100% async/await API (no callbacks)
  • Written natively in TypeScript
  • Watch mode for development (experimental - iterate your jobs without restarting worker)
  • If you're running really lean, you can run Graphile Worker in the same Node process as your server to keep costs and devops complexity down.

⁠Status

Production ready (and used in production).

We're still enhancing/iterating the library rapidly, hence the 0.x numbering; updating to a new "minor" version (0.y) may require some small code modifications, particularly to TypeScript type names; these are documented in the changelog.

This specific codebase is fairly young, but it's based on years of implementing similar job queues for Postgres.

To give feedback please raise an issue or reach out on discord: http://discord.gg/graphile⁠

⁠Requirements

PostgreSQL 10+* and Node 10+*.

If your database doesn't already include the pgcrypto extension we'll automatically install it into the public schema for you. If the extension is installed in a different schema (unlikely) you may face issues. Making alias functions in the public schema, should solve this issue (see issue #43⁠ for an example).

* Might work with older versions, but has not been tested.

⁠Installation

yarn add graphile-worker
# or: npm install --save graphile-worker

⁠Running

graphile-worker manages its own database schema (graphile_worker). Just point graphile-worker at your database and we handle our own migrations:

npx graphile-worker -c "postgres:///my_db"

(npx looks for the graphile-worker binary locally; it's often better to use the "scripts" entry in package.json instead.)

The following CLI options are available:

Options:
      --help                    Show help                              [boolean]
      --version                 Show version number                    [boolean]
  -c, --connection              Database connection string, defaults to the
                                'DATABASE_URL' envvar                   [string]
  -s, --schema                  The database schema in which Graphile Worker is
                                (to be) located
                                           [string] [default: "graphile_worker"]
      --schema-only             Just install (or update) the database schema,
                                then exit             [boolean] [default: false]
      --once                    Run until there are no runnable jobs left, then
                                exit                  [boolean] [default: false]
  -w, --watch                   [EXPERIMENTAL] Watch task files for changes,
                                automatically reloading the task code without
                                restarting worker     [boolean] [default: false]
      --crontab                 override path to crontab file           [string]
  -j, --jobs                    number of jobs to run concurrently
                                                           [number] [default: 1]
  -m, --max-pool-size           maximum size of the PostgreSQL pool
                                                          [number] [default: 10]
      --poll-interval           how long to wait between polling for jobs in
                                milliseconds (for jobs scheduled in the
                                future/retries)         [number] [default: 2000]
      --no-prepared-statements  set this flag if you want to disable prepared
                                statements, e.g. for compatibility with
                                pgBouncer             [boolean] [default: false]

⁠Library usage: running jobs

graphile-worker can be used as a library inside your Node.js application. There are two main use cases for this: running jobs, and queueing jobs. Here are the APIs for running jobs.

⁠run(options: RunnerOptions): Promise<Runner>

Runs until either stopped by a signal event like SIGINT or by calling the stop() method on the resolved object.

The resolved 'Runner' object has a number of helpers on it, see Runner object⁠ for more information.

⁠runOnce(options: RunnerOptions): Promise<void>

Equivalent to running the CLI with the --once flag. The function will run until there are no runnable jobs left, and then resolve.

⁠runMigrations(options: RunnerOptions): Promise<void>

Equivalent to running the CLI with the --schema-only option. Runs the migrations and then resolves.

⁠RunnerOptions

The following options for these methods are available.

  • concurrency: The equivalent of the CLI --jobs option with the same default value.
  • nohandleSignals: If set true, we won't install signal handlers and it'll be up to you to handle graceful shutdown of the worker if the process receives a signal.
  • pollInterval: The equivalent of the CLI --poll-interval option with the same default value.
  • logger: To change how log messages are output you may provide a custom logger; see Logger⁠ below
  • the database is identified through one of these options:
    • connectionString: A PostgreSQL connection string to the database containing the job queue, or
    • pgPool: A pg.Pool instance to use
  • the tasks to execute are identified through one of these options:
    • taskDirectory: A path string to a directory containing the task handlers.
    • taskList: An object with the task names as keys and a corresponding task handler functions as values
  • schema can be used to change the default graphile_worker schema to something else (equivalent to --schema on the CLI)
  • forbiddenFlags see Forbidden flags⁠ below
  • events: pass your own new EventEmitter() if you want to customize the options, get earlier events (before the runner object resolves), or want to get events from alternative Graphile Worker entrypoints.

Exactly one of either taskDirectory or taskList must be provided (except for runMigrations which doesn't require a task list).

One of these must be provided (in order of priority):

⁠Runner object

The run method above resolves to a 'Runner' object that has the following methods and properties:

  • stop(): Promise<void> - stops the runner from accepting new jobs, and returns a promise that resolves when all the in progress tasks (if any) are complete.
  • addJob: AddJobFunction - see addJob⁠.
  • promise: Promise<void> - a promise that resolves once the runner has completed.
  • events: WorkerEvents - a Node.js EventEmitter that exposes certain events within the runner (see WorkerEvents⁠).
⁠Example: adding a job with runner.addJob

See addJob⁠ for more details.

await runner.addJob("testTask", {
  thisIsThePayload: true,
});
⁠Example: listening to an event with runner.events

See WorkerEvents⁠ for more details.

runner.events.on("job:success", ({ worker, job }) => {
  console.log(`Hooray! Worker ${worker.workerId} completed job ${job.id}`);
});
⁠WorkerEvents

We support a large number of events via an EventEmitter. You can either retrieve the event emitter via the events property on the Runner object, or you can create your own event emitter and pass it to Graphile Worker via the WorkerOptions.events option (this is primarily useful for getting events from the other Graphile Worker entrypoints).

Details of what events we support and what data is available on the event payload is detailed below in TypeScript syntax:

export type WorkerEvents = TypedEventEmitter<{
  /**
   * When a worker pool is created
   */
  "pool:create": { workerPool: WorkerPool };

  /**
   * When a worker pool attempts to connect to PG ready to issue a LISTEN
   * statement
   */
  "pool:listen:connecting": { workerPool: WorkerPool };

  /**
   * When a worker pool starts listening for jobs via PG LISTEN
   */
  "pool:listen:success": { workerPool: WorkerPool; client: PoolClient };

  /**
   * When a worker pool faces an error on their PG LISTEN client
   */
  "pool:listen:error": {
    workerPool: WorkerPool;
    error: any;
    client: PoolClient;
  };

  /**
   * When a worker pool is released
   */
  "pool:release": { pool: WorkerPool };

  /**
   * When a worker pool starts a graceful shutdown
   */
  "pool:gracefulShutdown": { pool: WorkerPool; message: string };

  /**
   * When a worker pool graceful shutdown throws an error
   */
  "pool:gracefulShutdown:error": { pool: WorkerPool; error: any };

  /**
   * When a worker is created
   */
  "worker:create": { worker: Worker; tasks: TaskList };

  /**
   * When a worker release is requested
   */
  "worker:release": { worker: Worker };

  /**
   * When a worker stops (normally after a release)
   */
  "worker:stop": { worker: Worker; error?: any };

  /**
   * When a worker is about to ask the database for a job to execute
   */
  "worker:getJob:start": { worker: Worker };

  /**
   * When a worker calls get_job but there are no available jobs
   */
  "worker:getJob:error": { worker: Worker; error: any };

  /**
   * When a worker calls get_job but there are no available jobs
   */
  "worker:getJob:empty": { worker: Worker };

  /**
   * When a worker is created
   */
  "worker:fatalError": { worker: Worker; error: any; jobError: any | null };

  /**
   * When a job is retrieved by get_job
   */
  "job:start": { worker: Worker; job: Job };

  /**
   * When a job completes successfully
   */
  "job:success": { worker: Worker; job: Job };

  /**
   * When a job throws an error
   */
  "job:error": { worker: Worker; job: Job; error: any };

  /**
   * When a job fails permanently (emitted after job:error when appropriate)
   */
  "job:failed": { worker: Worker; job: Job; error: any };

  /**
   * When a job has finished executing and the result (success or failure) has
   * been written back to the database
   */
  "job:complete": { worker: Worker; job: Job; error: any  };

  /**
   * When the runner is terminated by a signal
   */
  gracefulShutdown: { signal: Signal };

  /**
   * When the runner is stopped
   */
  stop: {};
}>;

⁠Library usage: queueing jobs

You can also use the graphile-worker library to queue jobs using one of the following APIs.

NOTE: although running the worker will automatically install its schema, the same is not true for queuing jobs. You must ensure that the worker database schema is installed before you attempt to enqueue a job; you can install the database schema into your database with the following command:

yarn graphile-worker -c "postgres:///my_db" --schema-only

Alternatively you can use the WorkerUtils migrate method:

await workerUtils.migrate();
⁠makeWorkerUtils(options: WorkerUtilsOptions): Promise<WorkerUtils>

Useful for adding jobs from within JavaScript in an efficient way.

Runnable example:

const { makeWorkerUtils } = require("graphile-worker");

async function main() {
  const workerUtils = await makeWorkerUtils({
    connectionString: "postgres:///my_db",
  });
  try {
    await workerUtils.migrate();

    await workerUtils.addJob(
      // Task identifier
      "calculate-life-meaning",

      // Payload
      { value: 42 },

      // Optionally, add further task spec details here
    );

    // await workerUtils.addJob(...);
    // await workerUtils.addJob(...);
    // await workerUtils.addJob(...);
  } finally {
    await workerUtils.release();
  }
}

main().catch((err) => {
  console.error(err);
  process.exit(1);
});

We recommend building one instance of WorkerUtils and sharing it as a singleton throughout your code.

⁠WorkerUtilsOptions
  • exactly one of these keys must be present to determine how to connect to the database:
    • connectionString: A PostgreSQL connection string to the database containing the job queue, or
    • pgPool: A pg.Pool instance to use
  • schema can be used to change the default graphile_worker schema to something else (equivalent to --schema on the CLI)
⁠WorkerUtils

A WorkerUtils instance has the following methods:

  • addJob(name: string, payload: JSON, spec: TaskSpec) - a method you can call to enqueue a job, see addJob⁠.
  • migrate() - a method you can call to update the graphile-worker database schema; returns a promise.
  • release() - call this to release the WorkerUtils instance. It's typically best to use WorkerUtils as a singleton, so you often won't need this, but it's useful for tests or processes where you want Node to exit cleanly when it's done.
⁠quickAddJob(options: WorkerUtilsOptions, ...addJobArgs): Promise<Job>

If you want to quickly add a job and you don't mind the cost of opening a DB connection pool and then cleaning it up right away for every job added, there's the quickAddJob convenience function. It takes the same options as makeWorkerUtils as the first argument; the remaining arguments are for addJob⁠.

NOTE: you are recommended to use makeWorkerUtils instead where possible, but in one-off scripts this convenience method may be enough.

Runnable example:

const { quickAddJob } = require("graphile-worker");

async function main() {
  await quickAddJob(
    // makeWorkerUtils options
    { connectionString: "postgres:///my_db" },

    // Task identifier
    "calculate-life-meaning",

    // Payload
    { value: 42 },

    // Optionally, add further task spec details here
  );
}

main().catch((err) => {
  console.error(err);
  process.exit(1);
});

⁠addJob

The addJob API exists in many places in graphile-worker, but all the instances have exactly the same call signature. The API is used to add a job to the queue for immediate or delayed execution. With jobKey and jobKeyMode it can also be used to replace existing jobs.

NOTE: quickAddJob is similar to addJob, but accepts an additional initial parameter describing how to connect to the database).

The addJob arguments are as follows:

  • identifier: the name of the task to be executed
  • payload: an optional JSON-compatible object to give the task more context on what it is doing
  • options: an optional object specifying:
    • queueName: the queue to run this task under
    • runAt: a Date to schedule this task to run in the future
    • maxAttempts: how many retries should this task get? (Default: 25)
    • jobKey: unique identifier for the job, used to replace, update or remove it later if needed (see Replacing, updating and removing jobs⁠); can be used for de-duplication (i.e. throttling or debouncing)
    • jobKeyMode: controls the behavior of jobKey when a matching job is found (see Replacing, updating and removing jobs⁠)

Example:

await addJob("task_2", { foo: "bar" });

Definitions:

export type AddJobFunction = (
  /**
   * The name of the task that will be executed for this job.
   */
  identifier: string,

  /**
   * The payload (typically a JSON object) that will be passed to the task executor.
   */
  payload?: any,

  /**
   * Additional details about how the job should be handled.
   */
  spec?: TaskSpec,
) => Promise<Job>;

export interface TaskSpec {
  /**
   * The queue to run this task under (only specify if you want jobs in this
   * queue to run serially). (Default: null)
   */
  queueName?: string;

  /**
   * A Date to schedule this task to run in the future. (Default: now)
   */
  runAt?: Date;

  /**
   * Jobs are executed in numerically ascending order of priority (jobs with a
   * numerically smaller priority are run first). (Default: 0)
   */
  priority?: number;

  /**
   * How many retries should this task get? (Default: 25)
   */
  maxAttempts?: number;

  /**
   * Unique identifier for the job, can be used to update or remove it later if
   * needed. (Default: null)
   */
  jobKey?: string;

  /**
   * Modifies the behavior of `jobKey`; when 'replace' all attributes will be
   * updated, when 'preserve_run_at' all attributes except 'run_at' will be
   * updated, when 'unsafe_dedupe' a new job will only be added if no existing
   * job (including locked jobs and permanently failed jobs) with matching job
   * key exists. (Default: 'replace')
   */
  jobKeyMode?: "replace" | "preserve_run_at" | "unsafe_dedupe";

  /**
   * Flags for the job, can be used to dynamically filter which jobs can and
   * cannot run at runtime. (Default: null)
   */
  flags?: string[];
}

⁠Logger

We use @graphile/logger⁠ as a log abstraction so that you can log to whatever logging facilities you like. By default this will log to console, and debug-level messages are not output unless you have the environmental variable GRAPHILE_LOGGER_DEBUG=1. You can override this by passing a custom logger.

It's recommended that your tasks always use the methods on helpers.logger for logging so that you can later route your messages to a different log store if you want to. There are 4 methods, one for each level of severity (error, warn, info, debug), and each accept a string as the first argument and optionally an arbitrary object as the second argument:

  • helpers.logger.error(message: string, meta?: LogMeta)
  • helpers.logger.warn(message: string, meta?: LogMeta)
  • helpers.logger.info(message: string, meta?: LogMeta)
  • helpers.logger.debug(message: string, meta?: LogMeta)

You may customise where log messages from graphile-worker (and your tasks) go by supplying a custom Logger instance using your own logFactory.

const { Logger, run } = require("graphile-worker");

/* Replace this function with your own implementation */
function logFactory(scope) {
  return (level, message, meta) => {
    cons

Tag summary

Content type

Image

Digest

sha256:37bd9ad7e…

Size

53 MB

Last updated

4 months ago

docker pull graphile/worker