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()。

三种并发模式、共享代理池和 pause/resume

整数是固定并发,创建时至少为 1;运行后可以把目标改成 0..2**31-1:

worker = queue.worker(handle, concurrency=8)
run_task = asyncio.create_task(worker.run())
await worker.set_concurrency(0)  # pause:不再派发,已有 handler 不取消
await worker.set_concurrency(8)  # resume/resize

ResourceConcurrency 是同一进程/事件循环可由多个 queue/client 共享的总容量。容量是 admission permit,不是每个 Worker 的副本;缩容不会抢占已运行的 handler。

from ringo_task_queue import ResourceConcurrency

proxies = ResourceConcurrency(capacity=12)
west = client.queue("west").worker(handle, concurrency=proxies, weight=1)
east = client.queue("east").worker(handle, concurrency=proxies, weight=3)
runs = [asyncio.create_task(w.run()) for w in (west, east)]

await proxies.set_capacity(20)  # 例如健康代理数变化
# 共享 controller 时不能调用 west.set_concurrency(...)

await west.close()
await east.close()
await asyncio.gather(*runs)

持续有任务时,权重按平滑加权轮转长期接近 1:3;空闲成员的临时 probe 会在 250 ms 后撤回,因此容量会给活跃成员 work-conserving 使用。可手动暂停整组:await proxies.set_capacity(0),恢复时再设置非零容量。

AdaptiveConcurrency 使用 AIMD:当前容量个成功 ACK(且当时已饱和)后 +1,只有成功 NACK 且 RetryTaskError(congestion=True) 才按 *0.5 降低,始终限制在 minimum..maximum;普通异常、永久失败和 lease/传输丢失不改变它:

from ringo_task_queue import AdaptiveConcurrency, RetryTaskError

limits = AdaptiveConcurrency(initial=4, minimum=1, maximum=32)

async def handle(task, ctx):
    try:
        await call_upstream(task.payload)
    except RateLimitError as exc:
        raise RetryTaskError(str(exc), delay=2, congestion=True) from exc

worker = queue.worker(handle, concurrency=limits)

动态 controller 和运行中的 fixed resize 需要协议 minor 2;不改变 fixed 并发的 Worker 仍兼容旧 minor。controller 可以跨 queue/client 共享,但不能跨进程共享。

Handler 映射:

结果 动作
正常返回 ACK / SUCCEEDED
普通异常 NACK / 按策略重试
RetryTaskError NACK,可覆盖本次延迟;仅 congestion=True 且 NACK 被服务端接受时触发 Adaptive 的降速
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.2
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 已上传;安装开发版本时使用 --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。

许可证

MIT License。

Release files for ringo-task-queue 0.1.0.dev4

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.dev4
File
ringo_task_queue-0.1.0.dev4-py3-none-win_amd64.whl Python 3 none Windows x86-64 Details
ringo_task_queue-0.1.0.dev4-py3-none-manylinux_2_17_x86_64.whl Python 3 none Linux glibc 2.17+ x86-64 Details
ringo_task_queue-0.1.0.dev4-py3-none-manylinux_2_17_aarch64.whl Python 3 none Linux glibc 2.17+ ARM64 Details
ringo_task_queue-0.1.0.dev4-py3-none-macosx_11_0_x86_64.whl Python 3 none macOS 11.0+ x86-64 Details
ringo_task_queue-0.1.0.dev4-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.dev4-py3-none-win_amd64.whl

Download URL ringo_task_queue-0.1.0.dev4-py3-none-win_amd64.whl
Size 8.7 MB
Tags Python 3 Windows x86-64
SHA-256 checksum
How to use checksums
645a0e6d4e4e6859c88b4bae703e3eb3da58b0c0a69acdc249830dd632912424
BLAKE2b-256 checksum
How to use checksums
8131e8b517332945739ce9803a7155fecc74da5e3fd1a287ffc253ed05541aa7
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.dev4-py3-none-manylinux_2_17_x86_64.whl

Download URL ringo_task_queue-0.1.0.dev4-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
0ee688a115feba0ca84e52c8a05f3f4a8e16943390544de92138ff3ed39d1f57
BLAKE2b-256 checksum
How to use checksums
507c13c1c41bb96d0d15cb0048388b3d06b09d9e77a528f23d48fe6735aa6c98
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.dev4-py3-none-manylinux_2_17_aarch64.whl

Download URL ringo_task_queue-0.1.0.dev4-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
92886501a11725ed6e7508c4f32f4bc78f880aecc66e9942e0da9a9c5d1b0698
BLAKE2b-256 checksum
How to use checksums
dc99d6302eede2c848c5165aecca49eb14e79b58cae2c9e16b3a5517d46ff385
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.dev4-py3-none-macosx_11_0_x86_64.whl

Download URL ringo_task_queue-0.1.0.dev4-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
eb39265b6b227f89b5a9d787680e252556fbe5b1f1f8d6411fa75595b6fbccc0
BLAKE2b-256 checksum
How to use checksums
5bd70bb416495f184a29a1e6e74be4960a4bc4d6f1e4b8bce0215b4b1ed6a3db
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.dev4-py3-none-macosx_11_0_arm64.whl

Download URL ringo_task_queue-0.1.0.dev4-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
f33aae00135236d02365ccbc3d4d0785d26a240aa2d7194d75eff36ca8aebf8b
BLAKE2b-256 checksum
How to use checksums
8a5dc650ad909def2cbc5a4ffd44b62fdfe6733f5391a31348702c27dbfd6311
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