在 DBOS 示例中比较公平排队、限流与防抖

原文:Advanced Queue Patterns。作者归属:DBOS, Inc. / DBOS 文档与示例贡献者。中文翻译与核验:未完纪 / Codex。原页未明确注明首发日期;本文核验日期为 2026-10-05。

DBOS:并发、启动速率与最后输入:公平排队:租户分区 → 并发管理器 → 工作队列,管理器等待子工作结束;源码每租户最多 2,单 worker 最多 4;速率限制:工作队列 → 启动窗口,全局每 10 秒最多启动 2 个;这不是同时运行数量上限;防抖:同一租户不断提交 → 重置等待窗口,源码停止输入 10 秒后入队,使用最后一次输入;文档与源码要分开验收,原文示例为 5 / 1 / 5 秒;核验的运行参数为 4 / 2 / 10 秒
图:DBOS:并发、启动速率与最后输入。根据原文与静态核验内容自主绘制,非软件界面截图。

三种队列模式解决三种问题

这个示例演示如何用 DBOS 构建几种进阶队列模式。完整队列文档见 队列教程,全部示例源码见 官方 GitHub 示例目录。下面完整翻译原文的公平排队、速率限制、防抖和运行说明,并合并 README 中必需的启动步骤。

公平排队

假设一个队列的容量有限,需要公平分配给多个租户。应用部署在有限数量的服务器上,每台服务器只能同时处理 5 个工作流。你不希望某个租户独占所有容量,于是限制每个租户同一时刻最多运行 1 个工作流,同时仍允许每台服务器最多运行 5 个。

在 DBOS 中,可以把分区队列与普通的非分区队列组合起来:分区队列负责每租户限制,非分区队列负责每个服务器的限制。先注册两个队列并定义实际工作流:

DBOS.register_queue("concurrency-queue", worker_concurrency=5)
DBOS.register_queue("partitioned-queue", partition_concurrency=1)

# 公平排队:最多同时运行五个,但每个租户最多一个
@DBOS.workflow()
def fair_queue_workflow():
    time.sleep(5)

DBOS.register_queue 会把队列配置写入系统数据库,因此必须在 DBOS.launch() 之后调用。接着创建入队接口。它不直接把实际工作流入队,而是把一个“并发管理器”工作流提交到分区队列,以实施每租户限制:

@api.post("/workflows/fair_queue")
def submit_fair_queue(tenant_id: str):
    # 让分区队列限制每个分区的并发数
    with SetEnqueueOptions(queue_partition_key=tenant_id):
        DBOS.enqueue_workflow("partitioned-queue", fair_queue_concurrency_manager)

并发管理器连接两个队列:它向非分区队列提交实际工作流,并等待该工作流完成。

@DBOS.workflow()
def fair_queue_concurrency_manager():
    # 等待实际工作流结束,让两层流量限制同时生效
    return DBOS.enqueue_workflow("concurrency-queue", fair_queue_workflow).get_result()

管理器与实际工作流具有相同的生命周期,因此在实际工作结束前,租户对应的管理器始终占用分区并发名额。这个模式同时遵守分区队列的每租户限制,以及非分区队列的工作进程并发限制。也可以调整这一模式,将每租户限制与其他总体流量限制组合。

速率限制

有时需要限制一段时间内可以启动多少个工作流。调用具有速率配额的 API(例如许多大语言模型 API)时,这尤其有用。给队列配置速率限制即可实现:

DBOS.register_queue("rate-limited-queue", limiter={"limit": 2, "period": 10})

# 每 10 秒最多启动两个工作流
@DBOS.workflow()
def rate_limited_queue_workflow():
    time.sleep(5)

如果速率限制的 limit 为 X、period 为 Y,则每 Y 秒最多启动 X 个工作流。这个限制在使用同一队列的全部 DBOS 进程之间全局生效。入队方法与其他队列相同:

@api.post("/workflows/rate_limited_queue")
def submit_rate_limited_queue():
    DBOS.enqueue_workflow("rate-limited-queue", rate_limited_queue_workflow)

防抖

有时会在短时间内连续收到多个启动请求,但只希望实际启动一次。例如用户连续编辑输入框时,可以等最后一次编辑后经过一段时间,再启动处理工作流。

防抖会把执行推迟到距最近一次调用已经过指定时间之后。先定义工作流、队列以及相应的防抖器:

DBOS.register_queue("debouncer-queue")

@DBOS.workflow()
def debouncer_workflow(tenant_id: str, input: str):
    print(f"Executing debounced workflow for tenant {tenant_id} with input {input}")
    time.sleep(5)

debouncer = Debouncer.create(debouncer_workflow, queue="debouncer-queue")

随后通过防抖器提交工作流。对于每个租户,工作流会等到该租户最后一次提交输入后经过指定时间才入队;实际启动时,使用防抖器最后收到的输入。

# 每次提交都会更新这个租户的防抖窗口
# 文档示例:停止输入五秒后,把最后一次输入提交给工作流
@api.post("/workflows/debouncer")
def submit_debounced_workflow(tenant_id: str, input: str):
    debounce_key = tenant_id
    debounce_period_sec = 5
    debouncer.debounce(debounce_key, debounce_period_sec, tenant_id, input)

更多信息见 Debouncer 参考文档。

亲自运行示例

克隆并进入示例仓库:

git clone https://github.com/dbos-inc/dbos-demo-apps.git
cd dbos-demo-apps/python/queue-patterns

原文接下来要求按 README 运行。README 的完整设置步骤是先安装依赖,再启动应用:

uv sync
./launch_app.sh

然后打开 http://localhost:8000 查看队列演示。

版本核验:文档数字与源码数字

项目 在线文档 固定提交的运行源码
普通队列 worker_concurrency 5 4
分区队列 partition_concurrency 1 2
防抖等待时间 5 秒 10 秒
速率限制 10 秒最多启动 2 个 相同
模拟工作耗时 time.sleep(5) 相同

源码的关键配置为:

DBOS.launch()
DBOS.register_queue("concurrency-queue", worker_concurrency=4)
DBOS.register_queue("partitioned-queue", partition_concurrency=2)
DBOS.register_queue("rate-limited-queue", limiter={"limit": 2, "period": 10})
DBOS.register_queue("debouncer-queue")

# 防抖接口内的实际值
debounce_period_sec = 10

源码的注释仍写“五个并发、每租户一个”和“五秒防抖”,但运行参数应以真正执行的赋值与注册语句为准。应用还提供随机批量请求和阶段视图:随机批量接口创建 50 个请求,在前四个租户中随机选一个,给予它两倍抽样权重;五个租户中的 ed 没有参与这个随机批次。

公平排队视图汇总管理器的 ENQUEUED 与 PENDING 状态、普通队列正在运行的工作流,以及最近 30 分钟成功的管理器记录。它通过父工作流 ID 把普通队列中的任务关联回租户。防抖视图展示 DELAYED、PENDING 和最近 30 分钟的 SUCCESS;源码说明还存在等待入队的 ENQUEUED 阶段,但这个接口没有单独返回该阶段,因此界面不是所有瞬间状态的完整清单。

Python 项目要求 Python >=3.10,锁定 dbos==3.2.0。前端 package.json 声明 react ^19.2.0、vite ^7.2.4;本次读取的锁文件实际将 Vite 固定在 7.3.6,要求 Node ^20.19.0 || >=22.12.0。同一锁文件里的 ESLint 10.8.0 及 @eslint/js 10.0.1 要求更严格的 ^20.19.0 || ^22.13.0 || >=24;选择 Node 时应满足整套依赖的要求。

静态审查发现与演示边界

  • 源码用 uvicorn.run(app, host="0.0.0.0", port=8000) 监听全部网卡。接口没有身份认证或租户归属检查;外部可达时,调用者可以伪造 tenant_id、批量入队或查询工作流记录。建议本地演示将监听改为 127.0.0.1;生产接入时必须由已认证身份派生租户键,并对接口授权和限制请求。此建议是译者修正方案,不是原文已实施措施。
  • 防抖工作流将租户与输入原文打印到日志,状态接口也返回输入;示例中不应放入密码、密钥或真实敏感信息。
  • launch_app.sh 执行 npm install,会安装依赖并可能运行依赖生命周期脚本;应审查可信来源与锁文件。为了复现可采用经过审查的固定提交及 npm ci,这是附加建议,原脚本仍是 npm install。
  • 示例的 time.sleep(5) 只是模拟工作,不是实际 API 调用或生产副作用处理。把真实操作接入时,应按 DBOS 工作流与步骤语义处理恢复、重试和幂等性,不能由这段演示推断外部操作只会发生一次。
  • 本次只阅读网页、README、运行脚本、main.py、依赖声明与锁文件,未安装依赖、启动 Web 服务、创建数据库或运行负载测试。

来源、版本与许可说明

DBOS 页面版权为 DBOS, Inc.;dbos-demo-apps 树未见根 LICENSE。DBOS SDK 软件许可证不替代教程或仓库授权。

本译稿保留原文的技术步骤、示例与限制,并以“译注”或“补充核验”区分修正和新增说明。示例只做静态审查,未安装、运行或测试。

伴随源码核验提交:dbos-inc/dbos-demo-apps @ ddaf50e71fa8。在线文档可能与源码构建时间不同,文中已注明实质差异。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容