Skip to main content
Pre-release

This release is a pre-release and may not be stable for production use.

Ringo Task Queue — Python SDK

面向 Python 3.11+ 的 asyncio 任务队列 SDK。它支持 SDK 管理的本地 SQLite/PostgreSQL daemon 和远程 gRPC 服务,提供永久去重、至少一次执行、lease fencing、自动续租、重试和 DEAD/DLQ 语义。

安装后即可离线运行

pip install ringo-task-queue

正式发布包含五个同版本平台 wheel:

系统 wheel tag 内置 daemon
Windows x64 win_amd64 Windows x64
Linux x64 manylinux_2_17_x86_64 Linux x64
Linux ARM64 manylinux_2_17_aarch64 Linux ARM64
macOS Intel macosx_11_0_x86_64 macOS Intel
macOS Apple Silicon macosx_11_0_arm64 macOS ARM64

pip 会自动选择当前平台的 wheel。每个 wheel 已包含匹配的 Go daemon;本地模式不会在安装或运行时下载二进制,不要求安装 Go,也不需要设置 RINGO_DAEMON_PATH。PostgreSQL 服务本身不包含在包内。

开发版本需要允许预发布:

pip install --pre ringo-task-queue

快速开始

import asyncio
from ringo_task_queue import Ringo


async def main() -> None:
    async with Ringo.local("./data/ringo") as client:
        queue = client.queue("email")
        handled = asyncio.Event()

        async def handle(task, ctx) -> None:
            print(task.id, task.payload)
            await ctx.report_progress(1, 1, "sent")
            handled.set()

        worker = queue.worker(
            handle,
            task_types=["send-email"],
            concurrency=4,
        )
        worker_task = asyncio.create_task(worker.run())

        result = await queue.enqueue(
            unique_key="welcome:user-42",
            task_type="send-email",
            payload={"user_id": 42},
        )
        print(result.task_id, result.created)

        await asyncio.wait_for(handled.wait(), timeout=10)
        await worker.close()
        await worker_task


asyncio.run(main())

(queue, unique_key) 在任务完整生命周期内永久唯一。重复 enqueue 返回原任务 ID 和 created=False,不会创建副本,也不会抛出重复异常。

本地模式

SQLite

async with Ringo.local("./data/ringo") as client:
    queue = client.queue("crawl")

SDK 自动完成 daemon 校验、随机 loopback 端口启动、健康检查、协议协商和退出回收。SQLite 数据保存在 data_dir;daemon 日志追加到 data_dir/ringo-daemon.stderr.log。同一时间不能有两个本地 daemon 占用同一个 data dir。

PostgreSQL

import os
from ringo_task_queue import PostgresStorage, Ringo


async with Ringo.local(
    "./data/ringo-process",
    storage=PostgresStorage(os.environ["RINGO_POSTGRES_DSN"]),
) as client:
    queue = client.queue("crawl")

PostgreSQL 必须由应用提供。SDK 不下载、启动或停止 PostgreSQL。DSN 通过子进程环境传递,不放进进程参数,PostgresStorage 的 repr 也会隐藏它。

远程模式

async with Ringo.connect(
    "https://queue.internal.example:7233",
    token="service-token",
) as client:
    queue = client.queue("crawl")

grpc:///http:// 使用明文,grpcs:///https:// 使用 TLS;省略端口时默认 7233。

受限链路可以启用整个 gRPC channel 的 gzip:

async with Ringo.connect(
    "https://queue.internal.example:7233",
    token="service-token",
    compression="gzip",
) as client:
    ...

compression 只接受 "none" 或 "gzip"。它覆盖 unary RPC 和 Work stream 并在重连后保持;本地模式始终不压缩。

提交任务

from datetime import datetime, timedelta, timezone
from ringo_task_queue import RetryPolicy, RetryStrategy, TaskSpec

spec = TaskSpec(
    unique_key="page:10001",
    task_type="fetch-page",
    payload={"url": "https://example.com/10001"},
    priority=10,
    available_at=datetime.now(timezone.utc) + timedelta(seconds=5),
    retry_policy=RetryPolicy(
        max_attempts=5,
        strategy=RetryStrategy.EXPONENTIAL,
        initial_delay=timedelta(seconds=1),
        max_delay=timedelta(seconds=60),
        multiplier=2.0,
        jitter=0.2,
    ),
)
result = await queue.enqueue(spec)

也可以使用关键字参数:

result = await queue.enqueue(
    unique_key="page:10001",
    task_type="fetch-page",
    payload={"url": "https://example.com/10001"},
)

payload 必须显式提供;显式 None 表示 JSON null。允许 JSON 对象、数组和标量,但:

  • UTF-8 JSON 编码不超过 1 MiB;
  • 不允许非有限浮点数;
  • 整数必须位于 [-2**53, 2**53];
  • JSON 对象的键必须是字符串;
  • 实际使用建议保持在 16–64 KiB 内,大对象只传外部引用。

批量提交

batch = await queue.enqueue_many(
    [
        TaskSpec(
            unique_key=f"page:{i}",
            task_type="fetch-page",
            payload={"page": i},
        )
        for i in range(100)
    ]
)
print(batch.created_count, batch.duplicate_count)

enqueue_many 接受任意数量任务,SDK 自动按服务端单批上限 1,000 条切分并顺序发送,最终聚合为单个 BatchResult;batch.results 保留逐项结果。若任一批次失败,异常直接抛出,不返回部分结果。

查询和手工重试

task = await queue.get_by_id(result.task_id)
task = await queue.get_by_unique_key("page:10001")

# 始终按业务 unique_key 解释
reopened = await queue.retry("page:10001")

# 显式按服务端 task ID
reopened = await queue.retry_by_id(result.task_id)

状态包括 SCHEDULED、READY、LEASED、SUCCEEDED、DEAD 和 CANCELLED。Task 还提供 attempt 历史、最后错误、进度和当前 lease 摘要。手工 retry 保留历史和永久去重关系。

Worker(推荐)

处理队列任务推荐直接使用 Worker,而不是下面的手工 Claim:Worker 自动管理注册、credit、lease 续租、并发和优雅关闭。

from datetime import timedelta
from ringo_task_queue import PermanentTaskError, RetryTaskError


async def handle(task, ctx) -> None:
    await ctx.report_progress(1, 2, "started")
    if should_retry(task):
        raise RetryTaskError("rate limited", delay=30)  # 秒
    if invalid(task):
        raise PermanentTaskError("invalid input")
    await process(task)


worker = queue.worker(
    handle,
    task_types=["fetch-page"],
    concurrency=32,
    lease_duration=timedelta(seconds=60),
    grace_period=timedelta(seconds=30),
)
await worker.run()

run() 一直运行到其他协程调用 await worker.close(),并且只能调用一次。常驻消费时用 asyncio.create_task(worker.run()) 启动;已知有限批次(见下)直接 await worker.run()。

Handler 映射:

结果 动作
正常返回 ACK / SUCCEEDED
普通异常 NACK / 按策略重试
RetryTaskError NACK,可覆盖本次延迟
PermanentTaskError Reject / DEAD
取消 不发送终态,等待 lease 到期重投

Worker 持有 lease token,自动续租并严格控制并发。TaskContext 提供 attempt、max_attempts、worker_id、cancelled、wait_cancelled() 和 report_progress()。lease 丢失、stream 断开或强制关闭时,context 会被取消。

这是至少一次系统:handler 的外部副作用必须使用 unique_key 或任务 ID 实现幂等。

有限批次与 on_idle(PyPI 推荐范式)

已知有限批次(例如"把这批 seed 跑完就退出")的推荐写法:先在 run() 之前提交幂等 seed,再给 Worker 配置一次性的 on_idle 回调,并直接 await worker.run():

import asyncio
import sys

from ringo_task_queue import Ringo


async def main() -> None:
    async with Ringo.local("./data/ringo") as client:
        queue = client.queue("email")

        # 1. 先提交幂等 seed(run() 之前 enqueue)。
        for user_id in (101, 102, 103):
            await queue.enqueue(
                unique_key=f"welcome:{user_id}",
                task_type="send-email",
                payload={"user_id": user_id},
            )

        async def handle(task, ctx) -> None:
            await send_welcome(task.payload["user_id"])

        # 2. 队列在该 Worker 的 task type 范围内连续空闲满 1 秒后触发一次。
        worker = queue.worker(
            handle,
            task_types=["send-email"],
            concurrency=4,
            on_idle=lambda: sys.exit(0),
        )

        # 3. 前台直接 await run();on_idle 触发后整个批次随之结束。
        await worker.run()


asyncio.run(main())

语义与限制(务必阅读):

  • 固定参数:SDK 以 250 ms 间隔轮询、连续空闲满 1 秒触发一次;这两个值是协议级常量,不提供配置参数。
  • 单次触发:回调最多触发一次。正常返回后 Worker 继续常驻运行且不再触发;要结束进程请在回调里退出(如 sys.exit(0)),异常会原样从 await worker.run() 传播。
  • type scope:活动数只统计该 Worker task_types 范围内的任务(SCHEDULED/READY/LEASED 都算活动),不影响队列里其他类型的任务。
  • 开放队列竞态:空闲观察只是开放队列上的观测快照,不是持久完成判定。回调正常返回且 Worker 保持常驻时,观察之后提交的新任务仍会被正常消费;但如果回调直接退出进程(如 sys.exit(0)),观测之后、退出之前竞态提交的任务不会被本 Worker 处理,可能留给其他 Worker 或下一次运行。因此 sys.exit(0) 范式只适合生产者已经停止的有限批次。需要严格分布式完成判定时,必须另行设计 queue seal/generation。
  • 不要在回调里 close():on_idle 由 worker.run() 协程内联执行,在回调里等待同一 Worker/client 的关闭会造成自等待。外层资源作用域(如上面的 async with)负责清理。
  • 兼容性:不带 on_idle 的 Worker 从不发送活动查询 RPC,与协议 minor 0 的旧 server 完全兼容;配置 on_idle 需要 minor ≥ 1 的 server,旧 server 会在 run() 快速抛出 IncompatibleVersionError。

手工 Claim

常规消费任务优先使用上面推荐的 Worker;手工 Claim 适用于需要自行控制 claim/ACK 时序的高级场景。

from datetime import timedelta

leased = await queue.claim(
    limit=10,
    task_types=["fetch-page"],
    wait_timeout=timedelta(seconds=5),
    lease_duration=timedelta(seconds=60),
    worker_id="manual-worker-1",
)

for item in leased:
    try:
        await process(item.task)
        await item.ack()
    except TimeoutError as exc:
        await item.nack(str(exc), delay=timedelta(seconds=10))
    except ValueError as exc:
        await item.reject(str(exc))

续租:

await item.extend_lease(timedelta(seconds=60))

lease token 不公开。旧 attempt 或迟到 ACK 会被 fencing 拒绝并映射成 LeaseLostError。

生命周期和关闭

推荐:

async with Ringo.local("./data/ringo") as client:
    ...

或手工:

client = Ringo.local("./data/ringo")
await client.start()
try:
    ...
finally:
    await client.close()

close() 可重复调用,且调用者取消不会中断底层清理。它会先 drain worker,在 grace period 内继续续租并允许 handler 完成,再关闭 channel 和 daemon。被强制取消的 handler 不发送 ACK/NACK。

本地 daemon 意外退出时最多自动重启三次。结果不确定的写操作不会被盲目重放;只读查询可以在重新协商后重试。

错误处理

from ringo_task_queue import NotFoundError, RingoError, UnavailableError

try:
    task = await queue.get_by_unique_key("missing")
except NotFoundError:
    ...
except UnavailableError as exc:
    logger.warning("temporarily unavailable: %s", exc)
except RingoError as exc:
    logger.exception(
        "queue error code=%s request_id=%s metadata=%r",
        exc.code,
        exc.request_id,
        exc.metadata,
    )

公共异常:

  • InvalidArgumentError;
  • DuplicateError(普通重复 enqueue 返回 created=False);
  • NotFoundError;
  • ConflictError;
  • LeaseLostError;
  • UnavailableError;
  • DeadlineExceededError;
  • IncompatibleVersionError;
  • InternalError。

默认值与限制

项目 默认值或限制
协议 1.1
schema 3
lease 60 秒
worker concurrency 1
graceful shutdown 30 秒
retry attempts 默认 3,最大 32
retry 指数退避,1 秒到 60 秒,multiplier 2,jitter 0.2
batch SDK 自动分批,单批最大 1,000
payload 最大 1 MiB UTF-8 JSON
error text 最大 256 KiB UTF-8
daemon 自动重启 最多 3 次

wheel 完整性和二进制解析

本地 daemon 解析优先级:

  1. Ringo.local(..., daemon_path=...);
  2. RINGO_DAEMON_PATH;
  3. wheel 内置平台 daemon;
  4. 仅用于开发 checkout 的搜索。

普通 PyPI 安装直接使用第 3 项。SDK 在启动前验证 daemon 是普通文件、长度和 SHA-256 与 manifest 一致,然后运行 version --json 校验产品名、版本、协议和 schema。校验失败时拒绝启动,不会静默使用错误二进制。

常见问题

No matching distribution found

确认 Python 为 3.11+,对应平台 wheel 已上传;安装 .dev0 时使用 --pre。

找不到 daemon

正式 wheel 不应出现。使用 python -m pip show -f ringo-task-queue 检查是否误装了 editable checkout 或非正式产物。

ConflictError

另一个进程通常已占用同一个 SQLite data dir。复用已有客户端或换一个目录。

任务重复执行

这是至少一次语义的正常边界。使用任务 ID/unique_key 为外部副作用建立幂等约束。

Worker 关闭较慢

默认允许 handler 在 30 秒内完成。可以调整 grace_period,并让 handler 响应 cancellation 或 ctx.cancelled。

维护者构建平台 wheel

真实 daemon 不提交到源码树。必须先从同一源码版本生成并验证五平台 release manifest,再逐个平台构建 wheel:

go run ./tools/release build --out dist --version 0.1.0-dev --commit <revision> --date <RFC3339-UTC>
go run ./tools/release verify --manifest dist/manifest.json

然后从 sdk/python 运行:

python tools/build_platform_wheel.py \
  --binary ../../dist/ringo-task-queue-linux-amd64 \
  --release-manifest ../../dist/manifest.json \
  --platform linux-amd64 \
  --output-dir ../../package-dist/python

分别构建 windows-amd64、linux-amd64、linux-arm64、macos-amd64、macos-arm64。构建会记录根 manifest 的 SHA-256 来源,完成后清理暂存二进制。上传前运行:

python tools/verify_wheel_set.py ../../package-dist/python
python -m twine check ../../package-dist/python/*

没有经过完整五目标 manifest 暂存的 wheel 和 sdist 会被主动拒绝,避免发布不含 daemon 或平台集合不完整的包。

许可证

MIT License。

Release files for ringo-task-queue 0.1.0.dev3

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

Built distributions (wheels)

Table of built distributions (wheels) for ringo-task-queue 0.1.0.dev3
File
ringo_task_queue-0.1.0.dev3-py3-none-win_amd64.whl Python 3 none Windows x86-64 Details
ringo_task_queue-0.1.0.dev3-py3-none-manylinux_2_17_x86_64.whl Python 3 none Linux glibc 2.17+ x86-64 Details
ringo_task_queue-0.1.0.dev3-py3-none-manylinux_2_17_aarch64.whl Python 3 none Linux glibc 2.17+ ARM64 Details
ringo_task_queue-0.1.0.dev3-py3-none-macosx_11_0_x86_64.whl Python 3 none macOS 11.0+ x86-64 Details
ringo_task_queue-0.1.0.dev3-py3-none-macosx_11_0_arm64.whl Python 3 none macOS 11.0+ ARM64 Details

Total release size: 41.5 MB

Release files / ringo_task_queue-0.1.0.dev3-py3-none-win_amd64.whl

Download URL ringo_task_queue-0.1.0.dev3-py3-none-win_amd64.whl
Size 8.7 MB
Tags Python 3 Windows x86-64
SHA-256 checksum
How to use checksums
d513bfaaf4acf840584003378c13b46d38e6a05a0707834c424af97a1b19df61
BLAKE2b-256 checksum
How to use checksums
01dc2837d56b1982fe932b366328d057b942e1029f495f2a242604d09399fa0b
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.11.9

Release files / ringo_task_queue-0.1.0.dev3-py3-none-manylinux_2_17_x86_64.whl

Download URL ringo_task_queue-0.1.0.dev3-py3-none-manylinux_2_17_x86_64.whl
Size 8.5 MB
Tags Linux glibc 2.17+ x86-64 Python 3
SHA-256 checksum
How to use checksums
95711189e05c3b1b2a3b5f0ca597a31ce0cc9d837f9386a23c3aee5c69d63b92
BLAKE2b-256 checksum
How to use checksums
03165fd7557dfa284845521d88de6536a30877bfc6f34f2c7547a35b82c3a33e
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.11.9

Release files / ringo_task_queue-0.1.0.dev3-py3-none-manylinux_2_17_aarch64.whl

Download URL ringo_task_queue-0.1.0.dev3-py3-none-manylinux_2_17_aarch64.whl
Size 7.7 MB
Tags Linux glibc 2.17+ ARM64 Python 3
SHA-256 checksum
How to use checksums
3a57abaf096f527c0e14483d62d63839c30f8c4cbef088b45caac5382f8661d2
BLAKE2b-256 checksum
How to use checksums
3ddd9f98b443d9abe7ea9993d7a4e8766b2fddb5164984cc09200228c805737c
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.11.9

Release files / ringo_task_queue-0.1.0.dev3-py3-none-macosx_11_0_x86_64.whl

Download URL ringo_task_queue-0.1.0.dev3-py3-none-macosx_11_0_x86_64.whl
Size 8.6 MB
Tags Python 3 macOS 11.0+ x86-64
SHA-256 checksum
How to use checksums
847098ec99fb0eb9617a233a0dd50425695d74773ca5d2d0f2ef5eb6b3142f4d
BLAKE2b-256 checksum
How to use checksums
a91a0be50791f748c3fdb93bb2be33ee15fc95628b8713bd3723d0c30b98fb1f
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.11.9

Release files / ringo_task_queue-0.1.0.dev3-py3-none-macosx_11_0_arm64.whl

Download URL ringo_task_queue-0.1.0.dev3-py3-none-macosx_11_0_arm64.whl
Size 8.0 MB
Tags Python 3 macOS 11.0+ ARM64
SHA-256 checksum
How to use checksums
2d6b400892914d2bbf0eafd6b666ce18d004fd6ca8c8230515705efb77b5c34b
BLAKE2b-256 checksum
How to use checksums
dfecefd4d4667236f30fd1e3acd479b5389c271275217584db56a0a88effb13e
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.11.9
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