ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Apache DolphinScheduler 任务执行框架深度解析:dolphinscheduler-task-executor 如何统一 run / track / report 任务实例

2026/9/15 12:53:16 拓冰建站 浏览量
Apache DolphinScheduler 任务执行框架深度解析:dolphinscheduler-task-executor 如何统一 run / track / report 任务实例 Apache DolphinScheduler 任务执行框架深度解析dolphinscheduler-task-executor 如何统一 run / track / report 任务实例【免费下载链接】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 的dolphinscheduler-task-executor模块——一个与具体任务类型无关的可复用任务执行框架。它定义了 Worker 端如何运行、如何跟踪、如何上报一个任务实例的完整机制任务如何被调度进容器、生命周期状态如何流转、生命周期事件如何通过进程内事件总线与跨节点 RPC 上报给 Master。读完本文你将理解 DolphinScheduler 中 Shell、Spark、Flink 等所有任务类型共享的执行底座掌握TaskEngine、执行容器、事件总线、状态机映射与远程上报器的内部实现并了解为该模块扩展新执行策略的 SPI 边界与必须遵守的工程约束。模块定位Worker 侧通用的任务执行底座在 DolphinScheduler 的整体架构中dolphinscheduler-task-executor源码根目录dolphinscheduler-task-executor是一个纯库plain library模块它被 Worker 进程dolphinscheduler-worker内嵌使用。它回答的问题只有一个任务被分发到 Worker 之后怎么跑起来、怎么感知状态变化、怎么把结果汇报给 Master。它的核心设计原则是与任务类型解耦通用执行机制调度、生命周期、事件、上报全部收敛在本模块内具体任务类型的行为Shell、Spark、SQL……由dolphinscheduler-task-plugin下的各个插件提供本模块通过插件接口与它们协作本模块的依赖极简从 pom.xml 可以看到它只依赖三样东西dolphinscheduler-eventbus—— 任务生命周期总线的底层基础设施dolphinscheduler-task-api—— 任务 DTO 与契约如TaskExecutionContextdolphinscheduler-common—— 通用工具线程工具、日志、JSON 等。主包路径为org.apache.dolphinscheduler.task.executor全部源码位于 src/main/java/org/apache/dolphinscheduler/task/executor。核心子包全景该模块按职责划分为 7 个子包构成了容器执行 → 事件驱动 → 生命周期监听 → 远程操作 → 数据载体 → 异常体系 → 工作线程的完整闭环子包职责代表性类型task.executor.container执行模型容器容器拥有任务生命周期ExclusiveThreadTaskExecutorContainer、SharedThreadTaskExecutorContainer、AbstractTaskExecutorContainertask.executor.eventbus进程内延迟事件总线驱动生命周期流转TaskExecutorEventBus、TaskExecutorEventBusCoordinatortask.executor.listener生命周期监听器start / finish / fail / timeout / killITaskExecutorLifecycleEventListener、TaskExecutorLifecycleEventListenertask.executor.operations经 RPC 传输的操作请求/响应dispatch、kill、pause、reassignTaskExecutorDispatchRequest、TaskExecutorKillRequest、TaskExecutorPauseRequest、TaskExecutorReassignMasterRequest及各自 Responsetask.executor.dto任务状态与执行上下文 DTOTaskExecutorDTOtask.executor.exceptions模块异常体系TaskExecutorRuntimeException、TaskExecutorNotFoundException等task.executor.worker支撑容器的工作线程实现TaskExecutorWorker、AbstractTaskExecutorWorker其中operations子包与dolphinscheduler-extract-worker中定义的线上传输类型是兄弟关系sibling——前者是模块内的请求/响应对象后者是实际走 RPC 的线格式二者一一对应。关键类型从门面到基础设施TaskEngine / ITaskEngineWorker 提交与控制任务的统一门面ITaskEngineITaskEngine.java定义了任务运行时引擎的对外契约继承AutoCloseable共 6 个方法public interface ITaskEngine extends AutoCloseable { void start(); // 启动任务引擎 ListTaskExecutorDTO queryTaskExecutors(); // 查询所有任务执行器信息 void submitTask(ITaskExecutor taskExecutor) // 提交任务可能抛 TaskExecutorRuntimeException throws TaskExecutorRuntimeException; void pauseTask(int taskExecutorId) // 暂停任务非阻塞仅发送暂停信号 throws TaskExecutorNotFoundException; void killTask(int taskExecutorId) // 杀死任务非阻塞仅发送 kill 信号 throws TaskExecutorNotFoundException; void close(); // 关闭任务引擎 }TaskEngineTaskEngine.java是该接口的唯一实现。它的submitTask展示了任务提交的完整调用链共四步分配容器taskExecutorContainerDelegator.getExecutorContainer()拿到容器注册executorContainer.dispatch(taskExecutor)将任务分配注册给某个具体 worker 线程随后taskExecutorRepository.put(taskExecutor)写入内存注册表发事件taskExecutor.getTaskExecutorEventBus().publish(TaskExecutorDispatchedLifecycleEvent.of(taskExecutor))发布DISPATCHED事件启动executorContainer.start(taskExecutor)真正把任务 fire 到执行线程上。值得注意的是pauseTask/killTask的注释明确写着This method will not block, only send the pause signal——它们只是向任务自己的事件总线发布TaskExecutorPauseLifecycleEvent/TaskExecutorKillLifecycleEvent真正的动作由事件监听器异步完成从而避免控制路径被任务执行阻塞。queryTaskExecutors则从注册表拉取全部任务映射为TaskExecutorDTO含任务实例 ID、任务名、任务类型、项目编码、工作流实例 ID/名称、开始时间。TaskEngineBuilder启动装配器TaskEngine不接受new直出而是通过TaskEngineBuilderTaskEngineBuilder.java装配四个核心依赖engineName、taskExecutorRepository、taskExecutorContainerDelegator、taskExecutorEventBusCoordinator。Worker 侧的装配实现在 PhysicalTaskEngineFactory.java它把 Spring 容器里的物理任务执行器仓储、容器提供者、事件总线协调器注入 Builder构建名为PhysicalTaskEngine的引擎——这也印证了本模块自身不依赖 Spring所有 Bean 的接线都在 Worker 侧完成。TaskExecutorRepository运行中任务的内存注册表TaskExecutorRepositoryTaskExecutorRepository.java是任务执行器的内存注册表接口为ITaskExecutorRepository。它是引擎、容器、事件总线协调器和远程上报器共享的任务索引submitTask时写入事件总线协调器靠它遍历所有任务执行器并 fire 各自的事件总线远程上报器靠它根据taskInstanceId找回ITaskExecutor从而拿到 Master 地址taskExecutionContext.getWorkflowInstanceHost()并触发FINALIZE事件。执行模型两种线程容器容器层定义了任务的执行模型即一个任务占用多少个执行线程。所有容器继承自AbstractTaskExecutorContainerAbstractTaskExecutorContainer.java它基于TaskExecutorContainerConfig构建// TaskExecutorContainerConfig 核心字段 private String containerName; Builder.Default private int taskExecutorThreadPoolSize Runtime.getRuntime().availableProcessors() * 2 1; // 默认CPU 核数 * 2 1容器构造时创建固定大小的 daemon 线程池命名格式容器名-worker-%d和等量的TaskExecutorWorker并立即把每个 worker 的start()提交进线程池。随后容器围绕TaskExecutorAssignmentTable任务→worker 的分配表提供四个操作dispatch找候选 worker →registerTaskExecutor注册→ 登记到分配表。若容器满载所有独占 worker 都忙直接抛TaskExecutorRuntimeException(All ExclusiveThreadTaskExecutorWorker are busy)start按分配表找到 workerfireTaskExecutor把任务点燃加入 active 集合并唤醒 worker 线程pause/kill转发给任务执行器自身finalize注销任务与分配表记录、发送远程日志、调用taskExecutor.finalizeTask()slotUsage返回活跃 worker 数占总 worker 数的比例供调度方评估容器负载。ExclusiveThreadTaskExecutorContainer一任务一线程ExclusiveThreadTaskExecutorContainer.java 的候选 worker 选择逻辑只有一行taskExecutorWorkers.selectIdleWorker()。它保证每个任务执行器独占一个 worker 线程适合 CPU/IO 密集、耗时长的重型任务如 Spark、Flink、Sqoop 等任务之间互不干扰也便于隔离日志与资源。代价是并发上限受 worker 数量限制——满载时新任务会被拒绝。SharedThreadTaskExecutorContainer线程池共享面向轻量任务SharedThreadTaskExecutorContainer.java 采用taskExecutorWorkers.roundRobinSelectWorker()轮询选择任务共享同一组 worker 线程适合轻量、短平快的任务如 HTTP、Shell 等。它提高了线程利用率但任务间共享执行线程需要任务插件本身对线程亲和性不敏感。选择哪种容器由ITaskExecutorContainerProviderSPI在 Worker 侧决定AbstractTaskExecutorContainer的注释还暴露了其内部的可测试性getTaskExecutorAssignmentTable()与getTaskExecutorWorkers()均标注VisibleForTesting供 Worker 集成测试直接断言分配行为。TaskExecutorWorker容器背后的执行循环TaskExecutorWorker.java 是容器的马夫。每个 worker 维护两类任务集合registeredTaskExecutors已注册分配给它但可能未启动的任务activeTaskExecutors已 fire正在被跟踪执行的任务。start()是一个永不退出的跟踪循环遍历所有 active 任务对尚未启动的调用taskExecutor.start()然后根据taskExecutor.getRemainingTrackDelay()见下节决定是立即trackTaskExecutorState还是等更近的跟踪时点当 active 集合为空时通过activeTaskExecutorEmptyCondition.await()挂起等待fireTaskExecutor用signalAll()唤醒。任何任务执行抛出的Throwable都会被捕获并走onTaskExecutorFailed失败通道保证单个任务的异常不会拖垮整个 worker 循环。任务执行器与状态机从 INITIALIZED 到终态AbstractTaskExecutor 的启动与跟踪流程AbstractTaskExecutorAbstractTaskExecutor.java是所有任务执行器的骨架基类构造时即把状态置为INITIALIZED。它的start()方法揭示了任务的执行流水线防重复启动检查startFlaginitializeTaskContext()设置taskExecutionContext的开始时间并以 JSON 打印上下文publishTaskRunningEvent()状态置为RUNNING并发布TaskExecutorStartedLifecycleEventRUNNING 事件initializeTaskPlugin()加载并初始化具体任务插件dryRun 短路若taskExecutionContext.getDryRun() Flag.YES直接跳过触发阶段、状态置为SUCCEEDED返回doTriggerTaskPlugin()真正触发任务插件执行。其中第 2、4、6 步对应TaskInstanceLogHeader.printInitializeTaskContextHeader()/printLoadTaskInstancePluginHeader()/printExecuteTaskHeader()三段日志分段方便按日志头定位任务运行阶段。状态跟踪由trackTaskExecutorState()承担若当前状态已是终态则直接返回否则更新latestStateTrackTime并调用抽象方法doTrackTaskPluginStatus()拉取插件状态再通过TaskExecutorStateMappings映射为统一状态。跟踪节流的核心是getRemainingTrackDelay()DEFAULT_TRACK_INTERVAL为 10 秒10_000msworker 循环据此计算距上次跟踪的剩余间隔避免对任务状态做无节制的轮询——这正是track环节的性能关键。状态枚举与表驱动映射任务执行器状态TaskExecutorStateTaskExecutorState.java共 6 个取值其中PAUSED、KILLED、FAILED、SUCCEEDED为终态INITIALIZED(false) → RUNNING(false) → PAUSED(true) / KILLED(true) / FAILED(true) / SUCCEEDED(true)插件侧返回的TaskExecutionStatus通过 TaskExecutorStateMappings.java 这张映射表转换为执行器状态插件状态TaskExecutionStatus执行器状态TaskExecutorStateRUNNING_EXECUTIONRUNNINGSUCCESSSUCCEEDEDFAILUREFAILEDKILLKILLEDPAUSEPAUSED其他/默认INITIALIZED映射表的默认分支返回INITIALIZED而非抛异常这意味着新增一个插件状态却忘记在映射表中登记会导致状态被静默吞掉见下文 Gotchas。事件驱动进程内延迟事件总线事件类型全集生命周期事件类型由 TaskExecutorLifecycleEventType.java 定义共 10 种DISPATCHED、RUNNING、RUNTIME_CONTEXT_CHANGE、PAUSE、PAUSED、KILL、KILLED、SUCCESS、FAILED、FINALIZE其中KILLED、PAUSED、FAILED、SUCCESS在isFinished()中标记为完成类事件——远程上报器正是依据它决定何时释放任务实例资源见下节。每个事件都实现ITaskExecutorLifecycleEvent并按可操作 / 可上报分成IOperableTaskExecutorLifecycleEvent与IReportableTaskExecutorLifecycleEvent两个侧面可操作事件如 PAUSE、KILL在本地驱动容器动作可上报事件如 RUNNING、SUCCESS、FAILED、KILLED则需要发往 Master。TaskExecutorEventBus 与协调器TaskExecutorEventBusTaskExecutorEventBus.java继承自dolphinscheduler-eventbus的AbstractDelayEventBus是每个任务执行器私有的进程内延迟事件队列publish 时会以 JSON 形式记录日志且通过TaskLogMarkers.excludeInTaskLog()排除在任务日志之外避免刷屏。TaskExecutorEventBusCoordinatorTaskExecutorEventBusCoordinator.java则是总调度台start()启动一个单线程调度器每 50msDEFAULT_FIRE_INTERVAL轮询所有任务执行器的事件总线事件处理由一个大小为availableProcessors的 worker 线程池并发执行并用firingTaskExecutorIds集合防止同一任务的事件被并发重复 firedoFireTaskExecutorEventBus()从总线poll()队头事件按switch (event.getType())分发到所有已注册的ITaskExecutorLifecycleEventListener的对应回调onTaskExecutorDispatchedLifecycleEvent、onTaskExecutorStartedLifecycleEvent、onTaskExecutorKillLifecycleEvent……。事件被 fire 一次即出队若监听器处理抛异常则记录 error 日志——这意味着事件不重放监听器必须幂等。可靠上报TaskExecutorLifecycleEventRemoteReporterTaskExecutorLifecycleEventRemoteReporterTaskExecutorLifecycleEventRemoteReporter.java是本模块与 Master 通信的信使它把 Worker 本地发生的生命周期事件通过 RPC 送达 Master是Master 视图下任务状态的唯一真相来源。它的可靠性设计非常值得借鉴按任务分通道每个任务实例一个ReportableTaskExecutorLifecycleEventChannelLinkedBlockingQueue事件按序入队顺序发送 ACK 确认发送线程循环 peek 队头事件只有当收到 Master 的 ACKreceiveTaskExecutorLifecycleEventACK后才会按事件类型remove出队Master 地址取自taskExecutionContext.getWorkflowInstanceHost()失败重试DEFAULT_TASK_EXECUTOR_EVENT_RETRY_INTERVAL为3 分钟。事件若从未发送过latestReportTime null或距上次发送超过 3 分钟就会被重发重发等待通过waitIfAnyTaskExecutionEventChannelRetryIntervalPassed()精确休眠到下一个可重发时点避免忙等完成类事件的延迟回收当某个通道收到完成类事件getType().isFinished()的 ACK 且通道已空时才发布TaskExecutorFinalizeLifecycleEvent触发任务执行器的finalize清理——任务执行器的生命周期被刻意延长到整个事件上报闭环结束之后防止任务实例在事件尚未送达 Master 前就被销毁主机漂移感知onWorkflowInstanceHostChanged(taskInstanceId)会把对应通道内所有事件的latestReportTime清空迫使它们立即重发到新的 Master 地址。Worker 侧通过 PhysicalTaskExecutorLifecycleEventReporter.java 继承该类并注入实际的 RPC clientITaskExecutorEventRemoteReporterClient完成对接。扩展点 / SPI为 Worker 暴露的接口全集本模块虽然是个无 Spring 的库但它通过一组纯接口把自己暴露给 Worker构成清晰的扩展边界全部定义于 src/main/java/org/apache/dolphinscheduler/task/executorSPI 接口语义ITaskExecutor单个任务执行器配合抽象基类AbstractTaskExecutor定制执行/跟踪逻辑ITaskExecutorContainer/ITaskExecutorContainerProvider执行容器及其提供者可自定义执行模型ITaskExecutorFactory任务执行器工厂ITaskExecutorEventBusCoordinator事件总线协调器ITaskExecutorLifecycleEventListener生命周期事件监听器start / finish / fail / timeout / kill 全回调ITaskExecutorRepository运行中任务注册表ITaskExecutorStateTracker任务状态跟踪器ITaskExecutorWorker容器工作线程配合AbstractTaskExecutorWorker扩展从AbstractTaskExecutor与AbstractTaskExecutorContainer的存在可以看出新执行策略的标准姿势是继承这两个抽象类覆盖doTrackTaskPluginStatus/initializeTaskPlugin/doTriggerTaskPlugin执行器侧与getTaskExecutorWorkerCandidate容器侧这几个抽象方法而不是改动引擎与事件基础设施。实战约束本模块的工程红线GotchasCLAUDE.md 中明确列出了维护与二次开发时必须遵守的约束这些是踩坑总结逐条说明如下模块内无 Spring本模块是纯库不出现Component/Configuration所有 Bean 接线都在dolphinscheduler-worker完成参见 PhysicalTaskEngineFactory.java。不要在本模块添加 Spring 依赖否则会破坏 Worker 的装配边界状态转换是表驱动的TaskExecutorStateMappings是插件状态与执行器状态的唯一切换点。新增状态或转换只改一处而不同步更新映射表事件会被静默丢弃映射默认分支返回INITIALIZED就是证据同样地新增TaskExecutorLifecycleEventType时必须同步更新协调器的switch分发与监听器接口没有src/main/resources本模块无资源目录配置属于宿主进程Worker不要在这里添加任何配置资源生命周期事件是 Master 视图的真相来源如果 Master 认为某个任务卡住了最可能的根因是本模块中的某个生命周期事件从未被发布例如事件总线未 publish、上报通道重试失败、ACK 未收回。排查方向应聚焦事件链路而非 Master 侧逻辑本模块没有自己的单元测试覆盖率来自dolphinscheduler-worker的集成测试AbstractTaskExecutorContainer中的VisibleForTesting方法就是为此服务的。改动本模块代码后必须在 Worker 侧跑集成测试验证单靠本模块编译通过不足以证明正确性。模块关系图五兄弟各司其职最后从依赖关系与数据流两个维度看本模块在整个项目中的位置dolphinscheduler-eventbus提供AbstractDelayEventBus基类是本模块进程内任务生命周期总线的底层dolphinscheduler-task-api提供任务 DTO 与契约TaskExecutionContext、TaskExecutionStatus等是本模块与任务插件沟通的协议层dolphinscheduler-common提供线程工具ThreadUtils、日志标记TaskLogMarkers、JSON 与远程日志等公共设施dolphinscheduler-worker本模块目前唯一的生产消费者负责 Bean 装配、容器提供、RPC client 注入并承载全部集成测试dolphinscheduler-master通过 RPC 接收本模块 reporter 上报的生命周期事件从而构建 Master 侧的任务状态视图。数据流的完整闭环可以概括为Master 派发 → Worker 的 TaskEngine.submitTask → 容器 dispatch/start → AbstractTaskExecutor 执行并跟踪状态 → 生命周期事件写入 TaskExecutorEventBus → Coordinator 每 50ms 分发 → 监听器驱动本地动作pause/kill 等与远程上报 → RemoteReporter 按通道、按序、带 ACK 与 3 分钟重试上报 Master → Master 确认后触发 FINALIZE 完成资源回收。理解这条链路也就掌握了 DolphinScheduler Worker 端任务执行的全部精髓。【免费下载链接】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),仅供参考