拓冰建站拓冰建站
首页 / 资讯中心 / 正文

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

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用task的tags参数给需要限流的任务打标签。例如多个查询任务共用一个只允许 10 个连接的数据库就把它们都打上database这个 tagfrom prefect import flow, task task(tags[database]) def query_database(query: str): # Simulate database work return fResults 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 tags limit prefect concurrency-limit inspect database # Delete a concurrency limit prefect concurrency-limit delete databasecreate的用法是prefect concurrency-limit create TAG CONCURRENCY_LIMITinspect支持--output选项输出 JSON 格式目前仅支持 jsonls支持--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( tagdatabase, concurrency_limit10 ) # Read current limit for a tag limit await client.read_concurrency_limit_by_tag(tagdatabase) print(limit) # View all concurrency limits limits await client.read_concurrency_limits(limit10, offset0) print(limits) # Delete a concurrency limit await client.delete_concurrency_limit_by_tag(tagdatabase) 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/可选分支Terraformresource prefect_concurrency_limit database_limit { tag database concurrency_limit 10 }验证上限是否生效配置完成后用 CLI 查看已创建的上限prefect concurrency-limit ls prefect concurrency-limit inspect databasels会列出所有并发上限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_SECONDS60两点注意必须设在 Prefect 服务端而不是客户端文档对此有明确说明。默认值在两份文档中不一致概念页写的是“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),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门