Skip to main content

pgjobq

A job queue built on top of Postgres.

Project status

Please do not use this for anything other than experimentation or inspiration. At some point I may decide to support this long term (at which point this warning will be removed), but until then this is just a playground subject to breaking changes (including breaking schema changes).

Purpose

Sometimes you have a Postgres database and need a queue. You could stand up more infrastructure (SQS, Redis, etc), or you could use your existing database. There are plenty of use cases for a persistent queue that do not require infinite scalability, snapshots or any of the other advanced features full fledged queues/event buses/job brokers have.

Features

  • Best effort at most once delivery (jobs are only delivered to one worker at a time)
  • Automatic redelivery of failed jobs (even if your process crashes)
  • Low latency delivery (near realtime, uses PostgreSQL's NOTIFY feature)
  • Low latency completion tracking (using NOTIFY)
  • Dead letter queuing
  • Job attributes and attribute filtering
  • Job dependencies (for processing DAG-like workflows or making jobs process FIFO)
  • Persistent scheduled jobs (scheduled in the database, not the client application)
  • Job cancellation (guaranteed for jobs in the queue and best effort for checked-out jobs)
  • Bulk sending and polling to support large workloads
  • Back pressure / bound queues
  • Fully typed async Python client (using asyncpg)
  • Exponential back off for retries
  • Telemetry hooks for sampling queries with EXPLAIN or integration with OpenTelemetry.

Possible features:

  • Reply-to queues and response handling

Examples

from contextlib import AsyncExitStack

import anyio
import asyncpg  # type: ignore
from pgjobq import create_queue, connect_to_queue, migrate_to_latest_version

async def main() -> None:

    async with AsyncExitStack() as stack:
        pool: asyncpg.Pool = await stack.enter_async_context(
            asyncpg.create_pool(  # type: ignore
                "postgres://postgres:postgres@localhost/postgres"
            )
        )
        await migrate_to_latest_version(pool)
        await create_queue("myq", pool)
        queue = await stack.enter_async_context(
            connect_to_queue("myq", pool)
        )
        async with anyio.create_task_group() as tg:

            async def worker() -> None:
                async with queue.receive() as msg_handle_rcv_stream:
                    # receive a single job
                    async with (await msg_handle_rcv_stream.receive()).acquire():
                        print("received")
                        # do some work
                        await anyio.sleep(1)
                        print("done processing")
                        print("acked")

            tg.start_soon(worker)
            tg.start_soon(worker)

            async with queue.send(b'{"foo":"bar"}') as completion_handle:
                print("sent")
                await completion_handle.wait()
                print("completed")
                tg.cancel_scope.cancel()


if __name__ == "__main__":
    anyio.run(main)
    # prints:
    # "sent"
    # "received"
    # "done processing"
    # "acked"
    # "completed"

Development

  1. Clone the repo
  2. Start a disposable PostgreSQL instance (e.g docker run -it -e POSTGRES_PASSWORD=postgres -p 5432:5432 postgres)
  3. Run make test

See this release on GitHub: v0.10.0

Release files for pgjobq 0.10.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for pgjobq 0.10.0
File Size Uploaded
pgjobq-0.10.0.tar.gz 19.2 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for pgjobq 0.10.0
File Interpreter ABI Platform
pgjobq-0.10.0-py3-none-any.whl Python 3 none any Details

Total release size: 43.3 kB

Release files / pgjobq-0.10.0.tar.gz

Download URL pgjobq-0.10.0.tar.gz
Size 19.2 kB
Tags Source
SHA-256 checksum
How to use checksums
1deb0cb74b76ab1822de0f3e0a8913f1b94505268c7b65f9e950b6045d12a2a2
BLAKE2b-256 checksum
How to use checksums
9f4c580141800f827f03946483c0cc761b0c5355163293caa857488e96202db3
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via poetry/1.2.2 CPython/3.11.0 Linux/5.15.0-1022-azure

Release files / pgjobq-0.10.0-py3-none-any.whl

Download URL pgjobq-0.10.0-py3-none-any.whl
Size 24.1 kB
Tags Python 3
SHA-256 checksum
How to use checksums
8c0421c0584dbe8ce05a6351507666c42e4bd70c22ea58121c9a4b7b7b2165c5
BLAKE2b-256 checksum
How to use checksums
56572b38980cf3a1777a041bc0eb6f670240711a604c724f8d0300cc4c800f03
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via poetry/1.2.2 CPython/3.11.0 Linux/5.15.0-1022-azure

Release history Release notifications | RSS feed

This release

0.10.0 This release

2 release files

0.9.1

2 release files

0.9.0

2 release files

0.5.0

2 release files

0.4.0

2 release files

0.2.0

2 release files

0.1.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page