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")

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

若同一台机器上、同一用户的多个 SDK 进程需要共享一个 SQLite 队列,可显式启用共享模式:

async with Ringo.local("./data/ringo", shared=True) as client:
    result = await client.queue("crawl").enqueue(
        unique_key="page:1", task_type="fetch", payload={"page": 1}
    )

第一个客户端自动启动唯一持锁的共享 daemon,后续客户端认证后附着;启动者退出不影响仍在工作的客户端。最后一个客户端离开后 daemon 延时退出;崩溃后存活客户端会重新发现并接入,不会重放结果不明的写请求。无需手动常驻服务。该模式仅支持本机 SQLite,不支持跨用户、网络共享目录或 PostgreSQL。默认私有模式保持不变。

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.3
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.dev6

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

Total release size: 41.6 MB

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

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

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

Download URL ringo_task_queue-0.1.0.dev6-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
a1721222b2d42f98dc32a6bf1f11dab08907046bb94dbf104e6687aa0d9e1a57
BLAKE2b-256 checksum
How to use checksums
7b240e718453005375fb4eeda355a8b9ef57e494a29d88a256f7526500656208
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.11

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

Download URL ringo_task_queue-0.1.0.dev6-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
f5ddf6e7fbb831d66aa0e34b051a81750cf2d0f7702ef678154da9b5083cea44
BLAKE2b-256 checksum
How to use checksums
8ec317952865ec20675e60e81a514443a766036e3dcc48f6e1457f2743e99d0f
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.11

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

Download URL ringo_task_queue-0.1.0.dev6-py3-none-macosx_11_0_x86_64.whl
Size 8.7 MB
Tags Python 3 macOS 11.0+ x86-64
SHA-256 checksum
How to use checksums
5e71e2c48722e5bb8f0d8450888b9478bd73d21c61e79deebed5e726c2942c26
BLAKE2b-256 checksum
How to use checksums
2dc02dd5924fe1fbacd4444e768eae7d80d434b90571c499f21218ba72f34db7
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.11

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

Download URL ringo_task_queue-0.1.0.dev6-py3-none-macosx_11_0_arm64.whl
Size 8.1 MB
Tags Python 3 macOS 11.0+ ARM64
SHA-256 checksum
How to use checksums
146b50b5d108af6ae021a1d6d08dcbfb5e01a68f4f448abec00d9efdb536415c
BLAKE2b-256 checksum
How to use checksums
6c8dfcb2dfcec92ba540ba104cf53eedd38c6b06b3c221f54ad55521a7487903
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.11
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