任务调度大数据后端前端【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址https://gitcode.com/gh_mirrors/do/dolphinscheduler点击查看免费下载任务组是 Apache DolphinScheduler 中用于限制任务实例并发数量的机制通过给任务实例分配并发名额来保护下游资源数据库、接口、Hadoop 集群等不被瞬时打满。本文以任务组管理为核心先讲解资源中心中的任务组配置与使用步骤再深入 Master 端TaskGroupCoordinator的实现源码剖析获取资源—排队等待—释放唤醒的完整闭环帮助你既会用、又懂其底层原理。任务组能解决什么问题在一个工作流密集编排的场景下可能会有大量任务实例在同一时刻被分发到 worker 上执行。如果这些任务都去访问同一个下游系统例如某张数据库表、某个第三方接口、某个计算集群瞬时并发过高会直接压垮下游。任务组的定位正是控制任务实例的并发从而控制对下游资源的压力。需要注意的是任务组对资源的限制是项目级别的和租户没有关系。也就是说任务组要么作用于某个具体项目要么作用于整个系统全部项目可见可用其并发控制粒度与租户无关。这与 Hadoop 集群自身的队列管控不同——集群队列是 Hadoop 侧的管控手段而任务组是 DolphinScheduler 调度侧的自有手段两者可以叠加使用。任务组配置新建任务组在 DolphinScheduler UI 中依次点击【资源中心】→【任务组管理】→【任务组配置】即可进入任务组管理页面然后点击新建任务组新建任务组时需要填写以下信息配置项说明是否必填【任务组名称】任务组在被使用时显示的名称用于在任务定义中引用必填【项目名称】任务组作用的项目非必选项如果不选择则整个系统所有项目均可使用该任务组选填【资源容量】允许任务实例并发的最大数量即该任务组最多同时有多少个任务实例在运行必填从源码看任务组的核心字段与界面一一对应。TaskGroup 实体映射数据库表t_ds_task_group包含name任务组名称projectCode作用的项目编码未绑定项目时全局可用groupSize资源容量最大并发数useSize当前已使用的并发数status任务组状态Flag.YES/NO决定该任务组是否被启用userId创建者。其中useSize就是资源池已占用大小调度时通过groupSize - useSize计算剩余可用并发名额。查看任务组队列在任务组配置页面中点击查看队列按钮即可查看任务组的使用信息队列信息来自数据库表t_ds_task_group_queue对应实体 TaskGroupQueue可以看到每个排队任务的任务名、所属工作流、组内优先级、是否强制启动forceStart、是否在队列中inQueue以及队列状态status。页面默认按update_time desc, id desc排序展示见 TaskGroupQueueMapper.xml。任务组的使用任务组仅适用于由 worker 执行的任务。像【switch】节点、【condition】节点、【sub_process】子流程等由 master 直接负责流转的节点类型不受任务组控制——因为它们并不真正占用 worker 的计算资源。下面以 shell 节点为例说明如何在任务定义中使用任务组在任务定义的高级配置或运行参数中只需要配置红色框内的两部分配置项说明【任务组名称】任务组配置页面中显示的任务组名称。这里只能看到该项目有权限的任务组新建任务组时选择了该项目以及作用在全局的任务组新建任务组时没有选择项目【组内优先级】当出现资源等待时优先级高的任务会最先被 master 分发给 worker 执行数值越大优先级越高从源码看任务定义中配置的任务组信息最终写入TaskInstance的taskGroupId与taskGroupPriority字段。master 端对所有待分发任务排序时BaseTaskExecuteRunnable#compareTo 采用的优先级比较链是工作流实例优先级processInstancePriority数值小者优先任务实例优先级taskInstancePriority数值小者优先任务组内优先级taskGroupPriority数值大者优先-taskGroupPriorityCompareResult实现倒序任务首次提交时间越早提交越优先。可见组内优先级是在工作流、任务两级优先级之后、提交时间之前的第三顺位排序依据用于决定同一任务组内谁先拿到名额。任务组的实现逻辑任务组的完整生命周期由 Master 进程内的 TaskGroupCoordinator 守护线程统一管理。它继承BaseDaemonThread在 Master 启动时被拉起每轮处理完逻辑后固定休眠 5 秒ThreadUtils.sleep(Constants.SLEEP_TIME_MILLIS * 5)并在处理前后获取、释放注册中心上的MASTER_TASK_GROUP_COORDINATOR_LOCK分布式锁保证多 Master 部署时只有一个实例在协调任务组名额。获取任务组资源TaskGroupCoordinator对外提供两个关键方法needAcquireTaskGroupSlot(TaskInstance)判断任务是否配置了任务组。判定条件是taskGroupId 0且任务组status Flag.YES启用状态。acquireTaskGroupSlot(TaskInstance)为任务申请任务组名额。Master 在分发任务时的整体逻辑是对应源码注释中的伪代码流程if (needAcquireTaskGroupSlot(taskInstance)) { taskGroupCoordinator.acquireTaskGroupSlot(taskInstance); return; // 立即停止本次分发 } // 未配置任务组正常抛给 worker 执行也就是说如果任务没有配置任务组则正常抛给 worker 运行如果配置了任务组则在抛给 worker 执行之前先检查资源池。acquireTaskGroupSlot并不直接占用名额而是往t_ds_task_group_queue写入一条状态为WAIT_QUEUE-1、inQueue YES、forceStart NO的排队记录随后任务就停留在SUBMITTED_SUCCESS状态等待被唤醒。当任务组资源池剩余大小groupSize - useSize不满足时任务就一直排队等待直到其他任务结束释放名额后被唤醒。释放与唤醒当获取到任务组资源的任务结束运行后Master 会调用releaseTaskGroupSlot(TaskInstance)释放资源。releaseTaskGroupSlot最终走releaseTaskGroupQueueSlot将排队记录置为inQueue NO、status RELEASE2。该方法是幂等的如果该名额已经释放再次调用不会产生副作用。释放并不是终点。TaskGroupCoordinator的每轮循环会执行以下四个步骤见run()方法amendTaskGroupUseSize()校准任务组useSize与真实占用数countUsingTaskGroupQueueByGroupId一致防止异常场景下名额泄漏amendTaskGroupQueueStatus()若排队记录对应的TaskInstance已不存在或已结束则将其对应的队列记录修正为RELEASE避免僵尸队列长期占坑dealWithForceStartTaskGroupQueue()处理被强制启动的排队任务——先通过 RPC 唤醒对应的等待任务实例再把队列记录置为RELEASE移出队列dealWithWaitingTaskGroupQueue()核心分配逻辑——找出还有剩余名额availableSize 0的任务组按priority desc取出至多availableSize条等待中的队列记录逐条执行数据库扣减名额acquireTaskGroupSlot→ RPC 唤醒等待任务实例 → 将队列状态更新为ACQUIRE_SUCCESS1。唤醒动作通过notifyWaitingTaskInstance完成它先校验任务实例状态为SUBMITTED_SUCCESS、工作流实例状态为RUNNING_EXECUTION且宿主host非空然后通过IWorkflowInstanceService#wakeupTaskInstance向承载该工作流实例的 Master 发送 RPC 请求。任务实例被唤醒后会再次进入分发流程去真正获取任务组资源并运行。任务组流程图下面这张图直观展示了从任务申请资源到任务结束释放资源、唤醒下一个任务的完整流转队列状态机小结排队记录的状态定义见枚举 TaskGroupQueueStatus共三态状态code含义WAIT_QUEUE-1等待队列中等待任务组名额ACQUIRE_SUCCESS1已成功获取名额任务正在运行RELEASE2已释放移出队列从源码结构还可以推断任务组队列 SQL 均按priority desc排序见 TaskGroupQueueMapper.xml因此同组内组内优先级数值越大的任务越先被dealWithWaitingTaskGroupQueue选中并唤醒这与 UI 上数值越大优先级越高的说明完全一致。实践要点与注意事项作用范围是项目级任务组要么绑定具体项目仅该项目任务可用要么不绑定全局可用与租户无关。只约束 worker 任务【switch】、【condition】、【sub_process】等由 master 执行的节点不受任务组控制控制并发时请评估这些节点之外的 worker 任务。资源容量要按下游承受能力设定groupSize即最大并发数建议结合下游系统的吞吐上限设置避免设得过大失去保护意义设得过小造成任务长时间排队。组内优先级只影响排队顺序它不改变任务总量只在名额不够、需要排队时决定谁先获得名额数值越大越优先。多 Master 部署天然安全TaskGroupCoordinator每轮通过注册中心分布式锁互斥执行且每 5 秒一轮巡检还会自动校准useSize与队列状态即使出现异常也能自愈。总结任务组是 DolphinScheduler 在调度侧提供的轻量级并发控制方案配置上只需在任务组管理中设定名称、项目与资源容量在任务定义中选择任务组并设置组内优先级实现上由 Master 的TaskGroupCoordinator守护线程统一协调通过数据库排队记录t_ds_task_group_queue与 RPC 唤醒机制完成获取—排队—释放—唤醒的闭环。理解了这套状态机与轮询机制你在实际项目中就能更精准地设计并发上限也能更从容地排查任务长时间不执行、名额不释放等排队类问题。赞分享任务调度大数据后端前端【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址https://gitcode.com/gh_mirrors/do/dolphinscheduler点击查看免费下载相关推荐Apache DolphinScheduler 任务组管理Task Group完全指南并发控制原理、配置实操与源码解析Apache DolphinScheduler 任务组管理Task Group完全指南并发控制原理、配置实操与源码解析 本文围绕 docs/docs/zh任务调度数据编排工作流自动化后端大数据Apache DolphinScheduler 任务组Task Group详解并发控制、队列调度与实现原理Apache DolphinScheduler 任务组Task Group详解并发控制、队列调度与实现原理 本指南以 Apache DolphinSche任务调度数据编排工作流自动化后端大数据Apache DolphinScheduler 资源中心完整指南本地/HDFS/S3/OSS 存储配置、文件管理与任务组并发控制Apache DolphinScheduler 资源中心完整指南本地/HDFS/S3/OSS 存储配置、文件管理与任务组并发控制 导读 资源中心Resour任务调度数据编排工作流自动化后端大数据创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考