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

Apache DolphinScheduler 任务组(Task Group)详解:并发控制、队列调度与实现原理

Apache DolphinScheduler 任务组Task Group详解并发控制、队列调度与实现原理【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler本指南以 Apache DolphinScheduler 的“任务组Task Group”功能为主题系统讲解任务组如何在同一项目维度上限制任务实例并发数、缓解下游资源如 Hadoop 集群、数据库连接池压力以及如何在创建任务定义时为任务绑定任务组并配置运行优先级。读完本文你将掌握任务组的创建、查看、绑定与优先级配置的完整操作流程并理解 Master 端任务组协调器TaskGroupCoordinator从“获取资源”到“释放唤醒”的底层调度原理。任务组的作用与适用场景任务组Task Group主要用于控制任务实例的并发数量其设计初衷是控制除调度器本身之外的其他资源压力——例如 Hadoop 集群的负载集群自身通常也有队列控制机制任务组可视为调度层之上的第二道闸门。在创建新任务定义时可以为任务配置对应的任务组并配置该任务在任务组内运行的优先级任务组的资源限制是项目级别的与租户Tenant无关——这意味着同一任务组下所有项目的任务实例共享同一个资源池跨租户同样生效用户只能查看被授权项目下的任务组只有对某个项目拥有写权限时才能创建或更新该项目下的任务组。从数据模型看任务组对应数据库表t_ds_task_group其核心字段见 TaskGroup.java包括字段含义name任务组名称project_code所属项目编码为 0 时表示全系统全局可用group_size资源池大小即允许的最大并发任务实例数use_size当前已占用的资源数status任务组状态启用/禁用user_id创建者创建任务组在 Web 界面中依次点击Resources - Task Group Management - Task Group option - Create Task Group即可进入创建弹窗。创建时需要填写以下信息任务组名称Task group name任务组被引用时展示的名称同一用户名下不允许重名。在 TaskGroupServiceImpl.createTaskGroup 中会先调用taskGroupMapper.queryByName做重名校验重名时抛出TASK_GROUP_NAME_EXSIT。项目名称Project name任务组生效的项目此字段为可选项。若未选择项目则该任务组对全系统的所有项目均可用此时project_code记为 0见 TaskGroupMapper.xml 中 queryTaskGroupPagingByProjectCode 的project_code in (projectCode, 0)查询。资源池大小Resource pool size允许同时运行的任务实例的最大数量。服务端会校验groupSize 0并抛出TASK_GROUP_SIZE_ERROR因此该值必须是大于 0 的正整数见 TaskGroupServiceImpl.java。创建时服务端还会执行requireProjectPerm(loginUser, projectCode, true)校验当前用户对所选项目是否拥有写权限见 TaskGroupServiceImpl.java与文档中“只有写权限才能创建/更新项目任务组”的描述一致。查看任务组队列在任务组列表页点击“查看任务组使用信息”按钮可以进入任务组队列Task Group Queue视图。队列视图展示每个已进入任务组的任务实例的详细排队与运行状态队列中的关键列包括项目名称、任务名称、流程实例、任务组名称、优先级、启动状态、是否入队In queue、任务状态、创建/更新时间等。这些信息来自数据库表t_ds_task_group_queue其字段定义见 TaskGroupQueue.java——其中priority记录该任务在任务组内的优先级forceStart标记是否为强制启动inQueue标记是否已进入队列占用资源status为队列状态。在任务定义中使用任务组注意任务组仅适用于由 Worker 执行的任务节点。由 Master 执行的节点类型——如switch节点、condition节点、sub_workflow子流程等——不受任务组控制。下面以 Shell 节点为例。在任务定义配置页面中只需要配置红框内的两个部分即可完成任务组绑定任务组名称Task group name在下拉框中选择任务组。这里只能看到两类任务组一是创建任务组时指定了本项目、且当前项目有权限访问的任务组二是创建任务组时未指定项目、全系统全局可用的任务组。优先级Priority当任务组内有任务排队等待资源时高优先级的任务会被 Master 优先分发到 Worker 执行。数值越大优先级越高。优先级的取值范围对应Priority枚举见 Priority.javaHIGHEST(0)、HIGH(1)、MEDIUM(2)、LOW(3)、LOWEST(4)其中 0 为最高优先级、4 为最低优先级。在任务组队列中该值越大越优先被调度与优先级枚举数值含义相反任务组内部按“数值大者优先”排序。此外Master 端还提供“强制启动forceStart”机制强制启动的任务组队列不需要等待空闲资源即可被唤醒运行。实现原理获取任务组资源Master 在分发任务时会先判断任务是否配置了任务组。判断逻辑位于 TaskGroupUtils.isUsingTaskGroup即检查TaskInstance.getTaskGroupId() 0未配置任务组任务被正常抛给 Worker 运行配置了任务组Master 先检查任务组资源池的剩余大小是否满足当前任务运行——若资源池减一后仍满足use_size group_size则继续分发若不满足则退出本次任务分发任务进入等待队列等其他任务完成并释放资源后被唤醒。资源占用通过一条带条件的原子 SQL 实现见 TaskGroupMapper.xml 的 acquireTaskGroupSlotupdate t_ds_task_group set use_size use_size 1 where id #{id} and use_size group_size这条 SQL 本质上是一种乐观锁只有当use_size group_size资源池未满时更新才会命中行影响行数 0从而保证并发环境下多个任务同时抢占资源池不会超卖。对应的服务端封装见 TaskGroupDaoImpl.acquireTaskGroupSlot。任务组槽位的获取流程在接口 ITaskGroupCoordinator 中定义任务进入SUBMITTED_SUCCESS状态后状态机基类 AbstractTaskStateAction 依次调用needAcquireTaskGroupSlot与acquireTaskGroupSlot——后者并不阻塞等待而是立即在t_ds_task_group_queue中插入一条WAIT_QUEUE状态的记录后返回随后任务暂停分发等待协调器在资源可用时将其唤醒见 TaskGroupCoordinator.acquireTaskGroupSlot。实现原理释放与唤醒当占用任务组资源的任务执行完毕时任务组资源会被释放。释放后 Master 会检查当前任务组内是否有任务在等待若存在等待任务则标记其中优先级最高的任务为其创建新的可执行事件事件中保存该任务的 ID该任务随后获得任务组资源并开始运行。具体的释放动作由releaseTaskGroupSlot完成它会删除该任务对应的所有TaskGroupQueue记录见 TaskGroupCoordinator.releaseTaskGroupSlot并通过use_size use_size - 1归还资源池。任务结束含失败、被 kill时状态机都会调用releaseTaskInstanceResourcesIfNeeded确保槽位被正确释放见 AbstractTaskStateAction.onFatalEvent。对于仍在WAIT_QUEUE中排队、尚未获得资源即被暂停Pause或终止Kill的任务releaseWaitingTaskGroupSlot会将其从等待队列移除见 TaskGroupCoordinator.releaseWaitingTaskGroupSlot并在 TaskSubmittedStateAction 的暂停/终止事件处理中被调用。任务组调度流程图下图完整描述了任务组从分发、排队到释放资源的整体流程流程要点如下任务分发Master 开始分发任务判断是否为 Master 节点执行switch、condition、sub_workflow等由 Master 执行的节点不进入任务组队列继续正常执行判断是否属于任务组队列TGQ属于则检查任务组队列是否已存在不存在则创建队列Create TGQ创建失败则等待下一轮重试判断任务是否需要在任务组内执行若任务已在任务组队列中且资源池为空、且该任务优先级最高则通过数据库乐观锁更新资源占用随后将任务发送给 Worker 执行释放资源Master 等待任务执行结束释放任务在任务组中占用的资源并唤醒队列中下一个最高优先级任务如此循环。任务组协调器TaskGroupCoordinator的调度循环任务组的资源分配与唤醒并不是由单个任务事件直接触发的而是由 Master 端一个常驻后台线程TaskGroupCoordinator统一驱动见 TaskGroupCoordinator.java。它启动时会先休眠 1 分钟确保旧实例迁移后遗留的任务组槽位已被释放唤醒操作本身幂等之后每轮循环执行四件事amendTaskGroupUseSize以t_ds_task_group_queue中实际ACQUIRE_SUCCESS且非强制启动的记录数为准校正任务组的use_size防止异常场景下资源数漂移源码 L135-L155amendTaskGroupQueueStatus清理任务实例已不存在或状态已结束的僵尸队列记录源码 L160-L207dealWithForceStartTaskGroupQueue处理“强制启动”的等待队列——直接通知对应任务实例运行并释放队列源码 L209-L280dealWithWaitingTaskGroupQueue查询所有use_size group_size的可用任务组按优先级取出等待队列中的任务逐个执行acquireTaskGroupSlotAndNotify源码 L282-L333。其中第 4 步是整个调度的核心它在一个数据库事务内先执行acquireTaskGroupSlot乐观锁占用资源池、再执行acquireTaskGroupQueue将队列状态置为ACQUIRE_SUCCESS全部成功后才通过 RPC 调用notifyWaitingTaskInstance唤醒目标任务源码 L335-L351。唤醒消息TaskGroupSlotAcquireSuccessNotifyRequest会被发送到该工作流实例所在 Master 的ITaskInstanceController从而触发任务继续分发。协调器每轮循环结束后休眠 5 秒再进入下一轮。任务组队列状态由枚举 TaskGroupQueueStatus 定义WAIT_QUEUE(-1)表示排队等待、ACQUIRE_SUCCESS(1)表示已成功获取资源、RELEASE(2)为已废弃的释放态。上述调度逻辑在 TaskGroupCoordinatorTest 中有对应的单元测试覆盖如空实例参数校验、needAcquireTaskGroupSlot判定等可结合测试用例进一步理解各接口的边界行为。小结任务组是项目级别的并发闸门通过资源池大小限制同时运行的任务实例数与租户无关可在创建任务组时指定项目或设为全系统全局可用。使用上只需两步在任务定义中绑定任务组并设置优先级注意switch、condition、sub_workflow等由 Master 执行的节点不受任务组控制。调度上由 Master 端协调器驱动任务配置了任务组后先进入WAIT_QUEUE资源充足时通过乐观锁 SQL 占用槽位并 RPC 唤醒任务任务结束后释放槽位并唤醒下一个最高优先级任务。排障可关注两个数据表t_ds_task_group资源池与t_ds_task_group_queue排队与运行状态结合任务组队列视图即可直观观察每个任务的优先级、入队状态与最终执行状态。相关文档与源码入口任务组官方文档、TaskGroupCoordinator 实现、TaskGroup 实体、TaskGroupMapper.xml。【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

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

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