Prefect 如何用 tag-based concurrency limits 限制并发任务运行?
2026/9/15 13:15:39 网站建设 项目流程

Prefect 如何用 tag-based concurrency limits 限制并发任务运行?

【免费下载链接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.项目地址: https://gitcode.com/GitHub_Trending/pr/prefect

当多个 flow 里的任务同时访问同一个共享资源——比如只允许 10 个连接的数据库,或有速率限制的外部 API——直接并发跑起来容易把资源打爆。Prefect 的 tag-based concurrency limits 就是针对这个场景:给任务打上 tag,再为这个 tag 设定一个上限,指定数量的任务处于Running状态后,其余带同一 tag 的任务会被延迟,直到有并发槽位释放。本文按“打 tag → 设置上限 → 验证”这条路径说明完整操作,适用前提是 Prefect 服务已在运行(延迟行为发生在服务端,下文会说明哪些配置必须设在服务端)。

tag-based 限制的工作机制

先把几个行为规则搞清楚,避免设置后行为与预期不符:

  • 限制只作用于任务:tag-based concurrency limits 是 Prefect 任务专用的并发控制,针对的是带指定 tag 的 task run。
  • 检查时机:每当一个 task run 尝试进入Running状态时,服务端会检查它的 tag 是否还有可用槽位。
  • 无限制即不限流:没有设置并发上限的 tag,其任务可以无限并发。
  • 多 tag 任务需要所有 tag 都有槽位:任务带多个 tag 时,只有每一个tag 都有可用并发位才会运行。
  • 上限设为 0 会立即中止:把某个 tag 的并发上限设为 0,会直接中止(abortion)带该 tag 的任务运行,而不是像正常限制那样延迟它们。
  • 与 global concurrency limits 的关系:自 Prefect 3.4.19 起,tag-based 限制由 global concurrency limits 实现。每创建一个 tag-based 上限,Prefect 会自动创建一条名为tag:{tag_name}的全局并发上限。这是实现细节,通常对用户透明,但你在 UI 或 API 响应里可能会看到这些tag:开头的全局上限——看到它们不是配置出错。

tag-based 与手动创建的 global concurrency limits 可以达成类似效果,但范围不同:global limits 可以作用于任何 Python 操作,tag-based limits 只针对 Prefect 任务。如果你只控任务,用 tag-based 更直接。

第一步:给任务打 tag

@tasktags参数给需要限流的任务打标签。例如多个查询任务共用一个只允许 10 个连接的数据库,就把它们都打上database这个 tag:

from prefect import flow, task @task(tags=["database"]) def query_database(query: str): # Simulate database work return f"Results for {query}" @flow def data_pipeline(): # These will respect the "database" tag limit query_database("SELECT * FROM users") query_database("SELECT * FROM orders") query_database("SELECT * FROM products")

如果任务需要同时受多个资源约束,传多个 tag 即可:

from prefect import task @task(tags=["database", "analytics"]) def complex_query(): # This task needs available slots in both "database" AND "analytics" limits return "complex results"

此时任务能运行,但databasetag 还没有上限,所以并发不受约束。

第二步:为 tag 设置并发上限

上限可以通过 CLI、Python client、REST API 或 Terraform 设置。以下以 CLI 为主路径。

用 CLI 设置(主路径)

# Set a limit of 10 for the "database" tag prefect concurrency-limit create database 10 # View all concurrency limits prefect concurrency-limit ls # View details about a specific tag's limit prefect concurrency-limit inspect database # Delete a concurrency limit prefect concurrency-limit delete database

create的用法是prefect concurrency-limit create TAG CONCURRENCY_LIMITinspect支持--output选项输出 JSON 格式(目前仅支持 json);ls支持--limit--offset分页和--output。CLI 还提供prefect concurrency-limit reset TAG,用于重置该 tag 上的并发槽位,完整的子命令说明见 prefect concurrency-limit 参考。

可选分支:Python client

在代码或脚本中管理上限时,用get_client()

import asyncio from prefect import get_client async def manage_concurrency_limits(): async with get_client() as client: # Set a concurrency limit of 10 on the "database" tag await client.create_concurrency_limit( tag="database", concurrency_limit=10 ) # Read current limit for a tag limit = await client.read_concurrency_limit_by_tag(tag="database") print(limit) # View all concurrency limits limits = await client.read_concurrency_limits(limit=10, offset=0) print(limits) # Delete a concurrency limit await client.delete_concurrency_limit_by_tag(tag="database") asyncio.run(manage_concurrency_limits())

可选分支:REST API

直接调用服务端 API(地址以你的服务实际地址为准,文档示例为本地服务):

# Create a concurrency limit curl -X POST "http://localhost:4200/api/concurrency_limits/" \ -H "Content-Type: application/json" \ -d '{"tag": "database", "concurrency_limit": 10}' # Get all concurrency limits curl "http://localhost:4200/api/concurrency_limits/"

可选分支:Terraform

resource "prefect_concurrency_limit" "database_limit" { tag = "database" concurrency_limit = 10 }

验证上限是否生效

配置完成后,用 CLI 查看已创建的上限:

prefect concurrency-limit ls prefect concurrency-limit inspect database

ls会列出所有并发上限,inspect database显示该 tag 上限的详情,能确认上限值就是你设置的那个(如 10)。

运行时验证看任务行为:上限 10 时,同一时刻最多 10 个databasetag 的任务处于Running,超出的部分不会失败,而是延迟进入Running状态等待槽位释放——概念页描述的行为是“delay the transition to aRunningstate”,即任务被推迟而不是报错。多 tag 任务则表现为:只有所有 tag 都还有可用槽位时才会运行。

调整被延迟任务的等待时间

任务因并发上限被延迟时,服务端会让客户端等待一段时间后再重试进入Running状态。这个等待时间由PREFECT_SERVER_TASKS_TAG_CONCURRENCY_SLOT_WAIT_SECONDS控制:

prefect config set PREFECT_SERVER_TASKS_TAG_CONCURRENCY_SLOT_WAIT_SECONDS=60

两点注意:

  1. 必须设在 Prefect 服务端,而不是客户端,文档对此有明确说明。
  2. 默认值在两份文档中不一致:概念页写的是“30 seconds(或该设置指定的值)”,而 settings 参考 中tag_concurrency_slot_wait_seconds的默认值是10(最小值 0,另支持环境名PREFECT_TASK_RUN_TAG_CONCURRENCY_SLOT_WAIT_SECONDS)。实际以你服务端生效的配置为准,如需固定行为就显式执行上面的config set命令。

与 global concurrency limits 的边界

  • tag-based limits 只能作用于带 tag 的 Prefect 任务;global concurrency limits 通过concurrencycontext manager 管理任何 Python 操作的槽位,两者都能限流数据库连接、API 调用等场景。
  • 自 Prefect 3.4.19 起 tag-based 上限在底层就是 global limits,所以 UI/API 中出现的tag:{tag_name}全局上限是同一份数据,不要重复创建。
  • 其他作用域的并发控制(work pool、work queue、deployment 级别的 flow run 上限)与 tag-based 限制互不替代,各管各的对象。

常见现象对照

  • 任务被延迟而不是失败:这是限制生效的正常表现,被延迟的任务按等待时间重试进入Running
  • 设置了 0 后任务直接中止:这是设计行为而非故障;要恢复执行,把上限改回正数或prefect concurrency-limit delete <tag>
  • 某 tag 槽位异常占满:可以用prefect concurrency-limit reset <tag>重置该 tag 的槽位。
  • 看到陌生的tag:xxx全局上限:是 3.4.19 起 tag-based 上限的底层实现,属正常现象。

更多机制说明可参考 tag-based concurrency limits 概念页 和 how-to 文档。

【免费下载链接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.项目地址: https://gitcode.com/GitHub_Trending/pr/prefect

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询