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 解析优先级:
Ringo.local(..., daemon_path=...);RINGO_DAEMON_PATH;- wheel 内置平台 daemon;
- 仅用于开发 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)
| File | Reset | |||
|---|---|---|---|---|
| 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
|