Skip to main content

A lightweight task queue, easy to integrate. / 一个轻量级任务队列,易于集成到现有业务系统。

Project description

barkueue

English | 简体中文

barkueue (bark-queue, or simply bark), is a task queue designed to embed into an existing business system, which is:

  • Uses a relational database to store business data
  • Has no event processor, or an event processor is overkill for the use case

It's NOT intended to:

  • Act as a high-performance and/or low-latency task queue processing every transaction
  • Replace the native external function or service broker of a DBMS

Usage

Basic Usage

import time

import barkueue as bark

# Use a list as an in-memory datasource, or connect to a database via SQLAlchemy.
arr: list[bark.Task] = [
    bark.Task("dog.bark", "Bluey"),
    bark.Task("dog.woof", "Bingo"),
    bark.Task("dog", "I'm *unicorse*, I like to eat children"),
]
ds = bark.datasource.ArrayDataSource(arr)

# 4 workers, persistence sync interval of 5s
app = bark.app([ds], max_workers=4, fetch_interval=5)


@app.handler("dog.#")
def dog_handler(app: bark.Application, task: bark.Task) -> None:
    if task.message.find("*unicorse*") != -1:
        # Let exceptions propagate — barkueue handles task status
        raise Exception(f"{task.message}")
    # Simulate business logic
    time.sleep(1)
    print(f"I'm {task.message}, a little puppy!")


@app.handler("dog.bark")
def bark_handler(app: bark.Application, task: bark.Task) -> None:
    # Simulate business logic
    time.sleep(3)
    print(f"{task.message} says: bark bark!")


@app.handler("dog.woof")
def woof_handler(app: bark.Application, task: bark.Task) -> None:
    # Simulate business logic
    time.sleep(2)
    print(f"{task.message} says: woof woof!")


app.run()

Built-in DataSources

barkueue ships with 2 datasources: an in-memory datasource and a SQLAlchemy-based ORM persistent datasource.

ArrayDataSource — in-memory datasource.

arr: list[bark.Task] = [bark.Task("order.paid", '{"id":1}')]
ds = bark.datasource.ArrayDataSource(arr)

arr serves as both input and output — push() writes status back to the Task objects in arr.

SqlAlchemyDataSource — persists to the barkueue_task table. barkueue does not depend on SQLAlchemy; install it separately as needed.

from sqlalchemy import create_engine
engine = create_engine("mssql+pyodbc://...")
ds = bark.datasource.SqlAlchemyDataSource(engine)

Table schema: id (int PK), topic, message, due, status. fetch() pulls rows where status IS NULL, push() performs a batch UPDATE.

Extending DataSource

DataSource is the abstraction layer between barkueue and external storage. It defines three methods:

Method Description
fetch() Populate self.tasks with unprocessed tasks; must set task.adapter = self
update_status(task, status) Buffer a status update — do not persist immediately
push() Flush buffered status updates to storage in a batch

DataSource is a Protocol. Any object with the above three methods and a tasks attribute satisfies the protocol — no explicit inheritance required.

Status Update Flow

When a Worker finishes a task, it calls task.update_status(status), which delegates to adapter.update_status(). For performance, update_status() only writes to an in-memory buffer (e.g. a dict); the actual persistence happens in push().

Each cycle, the DataSyncWorker operates in this order:

ds.push()   → flush buffered status updates from the previous cycle
ds.fetch()  → reload unprocessed tasks

Push-before-fetch ensures completed tasks are not re-fetched. Status updates lost due to a process crash before push() is a known trade-off — handlers that require idempotency should account for this.

Thread Safety

push() implementations should use an atomic-swap pattern: swap out the buffer dict under a lock, then perform I/O outside the lock. Worker threads only need the lock to write to the new dict and are never blocked by I/O.

Custom DataSource

Implement the four protocol members:

from barkueue.datasource.type import DataSource
from barkueue.task import Task

class MyDataSource(DataSource):
    tasks: list[Task]

    def __init__(self, ...) -> None:
        self.tasks = []
        self._updated: dict[str, int] = {}
        ...

    def fetch(self) -> None:
        """Pull tasks with status None, set adapter=self, append to tasks."""
        ...

    def update_status(self, task: Task, status: int) -> None:
        """Buffer the update — do not write to storage directly."""
        self._updated[task.id] = status

    def push(self) -> None:
        """Flush buffered updates in self._updated to storage."""
        ...

To-dos

  • fetch-consume event loop
  • fetch interval control
  • datasource diff, merge status update into data sync worker
  • in-memory datasource
  • retry task
  • event system (e.g. @app.init)

License

MIT

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

barkueue-0.0.1.tar.gz (39.3 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

barkueue-0.0.1-py3-none-any.whl (13.5 kB view details)

Uploaded Python 3

File details

Details for the file barkueue-0.0.1.tar.gz.

File metadata

  • Download URL: barkueue-0.0.1.tar.gz
  • Upload date:
  • Size: 39.3 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: uv/0.11.2 {"installer":{"name":"uv","version":"0.11.2","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Arch Linux","version":null,"id":null,"libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}

File hashes

Hashes for barkueue-0.0.1.tar.gz
Algorithm Hash digest
SHA256 71fee4766175228f68c1765308e78fe26d7855b5de906930992822c0b114373a
MD5 fc3f9f4caa4c29845aaf29d33d5a49da
BLAKE2b-256 bc328064a7ef5b2653e181845b28df75643bc046a89f1c149344d74ee06d053d

See more details on using hashes here.

File details

Details for the file barkueue-0.0.1-py3-none-any.whl.

File metadata

  • Download URL: barkueue-0.0.1-py3-none-any.whl
  • Upload date:
  • Size: 13.5 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: uv/0.11.2 {"installer":{"name":"uv","version":"0.11.2","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Arch Linux","version":null,"id":null,"libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}

File hashes

Hashes for barkueue-0.0.1-py3-none-any.whl
Algorithm Hash digest
SHA256 52386abe843045e1df9279547975e9727157e99fd7ae50cb688887107e2acf03
MD5 5443ddd9a8ad737c07b423096dfd3b9f
BLAKE2b-256 0041cf59e57504367646557e9bc10d9a74dc83159140528263004ed3b116f2a3

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page