资讯动态

Apache DolphinScheduler Dinky 任务插件:Flink 作业编排的创建、配置与源码级原理

发布时间:2026/9/15 12:14:41 来源:尧图企业网站定制
Apache DolphinScheduler Dinky 任务插件Flink 作业编排的创建、配置与源码级原理【免费下载链接】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 的Dinky 任务类型任务类型标识DINKY展开讲解如何在 DolphinScheduler 工作流中创建 Dinky 任务节点、配置Dinky 地址、Dinky 任务 ID、上线作业与自定义参数并通过 DinkyTask.java 等源码剖析 Worker 调用 Dinky Open API 的完整链路。读完本文你将掌握 Dinky 任务的完整配置方法、0.x/1.x 双版本兼容机制、参数传递方式与作业状态跟踪原理能够在实际项目中直接落地FlinkSQL / Flink Jar / SQL 一键编排的运维方案。任务概览Dinky 任务能做什么Dinky 是一个围绕 Flink 的一站式实时计算平台支持 FlinkSQL、Flink Jar 与 SQL 的开发、调试和运维。在 DolphinScheduler 中Dinky Task用于创建并执行 Dinky 类型任务从而把 Dinky 的作业管理能力纳入统一的工作流调度体系。当 Worker 执行 Dinky 任务时并不会直接在本机运行 Flink 作业而是通过调用Dinky Open API触发 Dinky 平台上的作业并持续跟踪作业实例状态直至结束。这意味着 DolphinScheduler 承担的是调度 状态追踪的角色真正的作业执行与资源管理由 Dinky 平台及其背后的 Flink 集群完成。从插件架构看Dinky 任务是一个独立的任务插件模块位于 dolphinscheduler-task-dinky。插件通过 DinkyTaskChannelFactory.java 以AutoService(TaskChannelFactory.class)方式注册其getName()返回任务类型标识DINKY即 UI 上任务类型列表中的名称。创建 Dinky 任务节点在 DolphinScheduler Web UI 中按以下步骤创建 Dinky 任务进入项目管理 → 项目名称 → 工作流定义点击创建工作流按钮进入 DAG 编辑页面从左侧工具栏将 Dinky 任务图标宽 15px 的小图标拖拽到画布中即完成节点创建在右侧弹出的节点配置面板中填写任务参数见下一节保存并发布工作流后即可按调度计划或手动触发执行。任务参数详解Dinky 任务除继承 DolphinScheduler 的默认任务参数如节点名称、运行标志、描述、任务优先级、Worker 分组、失败重试次数、失败重试间隔、超时告警等外还有以下专有参数参数描述Dinky 地址Dinky 服务的 URL例如http://localhost:8888。Worker 将基于该地址拼接 Open API 路径发起 HTTP 调用。Dinky 任务 IDDinky 作业对应的唯一 ID。在 Dinky 平台的任务管理中可以获取参见下文任务示例。上线作业指定当前 Dinky 作业是否上线。若开启则被提交的作业必须处于已发布状态且当前无对应的 Flink Job 实例在运行时才允许提交成功。自定义参数从 Dinky 1.0 开始支持传递自定义参数。目前仅支持IN类型输入不支持OUT类型输出。支持${param}语法获取 DolphinScheduler 的全局参数或局部动态参数。参数校验与默认值从源码 DinkyParameters.java 可以看到参数对象仅包含三个字段其中addressDinky 服务地址必填taskIdDinky 任务 ID必填online是否上线作业默认值为false。checkParameters()方法要求address与taskId均非空才算参数合法否则任务在init()阶段就会抛出DinkyTaskException(dinky task params is not valid)直接失败。因此创建节点时这两个字段是必须填写的。任务示例创建 Dinky 节点与获取任务 ID以下两张图展示了创建 Dinky 任务节点的完整过程。第一张图为 DolphinScheduler 的 DAG 编辑界面左侧选中任务类型DINKY右侧 Current node settings 配置面板中填入了 Dinky 地址示例值http://127.0.0.1:18888与 Dinky 任务 ID示例值27并开启了 Online task 上线开关第二张图展示了在 Dinky 平台中如何定位作业的任务 ID在 Dinky 平台的 Job Info 区域可以查看到作业的 Job ID示例值为1该 ID 即 DolphinScheduler 中Dinky 任务 ID参数的取值来源实操建议先在 Dinky 平台完成 FlinkSQL / Flink Jar 作业的开发与发布拿到任务 ID 后再回到 DolphinScheduler 创建工作流节点可避免因作业未发布导致上线作业提交失败。源码级原理Worker 如何触发与跟踪 Dinky 作业Dinky 任务的核心实现类是 DinkyTask.java它继承自AbstractRemoteTask表明该任务属于远程提交型任务Worker 不直接执行作业而是把作业提交给外部系统并轮询其结果。第一步版本探测决定调用协议在handle()阶段任务首先调用getDinkyVersion(address)请求{address}/openapi/version接口探测 Dinky 版本。由于 Dinky 0.x 与 1.x 的 Open API 协议不兼容插件据此分流版本号以0开头如 0.6.5→ 走V0提交与跟踪逻辑GET 请求 datas响应字段其他版本如 1.0.0→ 走V1逻辑POST JSON 请求 data响应字段。若版本接口返回异常则回退按0版本处理。从 DinkyTaskConstants.java 可以看出所有 API 路由均以/openapi/为前缀常量路由用途GET_VERSION/openapi/version探测 Dinky 版本SUBMIT_TASK/openapi/submitTask提交作业ONLINE_TASK/openapi/onLineTask上线作业0.x 版本SAVEPOINT_TASK/openapi/savepointTask取消作业savepoint 方式GET_JOB_INFO/openapi/getJobInstance查询作业实例状态第二步提交作业V0 版本通过 GET 请求携带id任务 ID参数。若onlinefalse调用submitTask若onlinetrue调用onLineTask。响应中读取datas.success与datas.jobInstanceId。V1 版本通过 POST JSON 请求submitTask请求体包含id、isOnline、variables三个字段响应读取success、data.jobInstanceId。关于上线作业的约束必须已发布且无对应 Flink Job 实例在运行正是由 Dinky 平台侧在submitTask/onLineTask处理时校验的DolphinScheduler 仅负责透传isOnline标记并解析返回结果。第三步轮询跟踪状态提交成功后任务进入trackApplicationStatus轮询阶段循环调用getJobInstance接口携带id jobInstanceId根据返回的status字段决定走向FINISHED按提交时status是否在线映射退出码记录 appId格式为address-taskId任务结束FAILED/CANCELED/UNKNOWN读取error信息并标记失败其他状态Thread.sleep(3000)后继续轮询即默认每 3 秒查询一次。一个值得注意的边界情况若提交成功但响应中没有jobInstanceId例如普通 SQL 任务瞬时完成任务会直接把address-taskId作为 appId 并以成功状态结束避免无限等待。第四步取消作业当工作流被终止时cancelApplication()会调用savepointTask接口以typecancel的方式对 Dinky 作业执行取消savepoint操作。从 DinkyTaskConstants.java 可见SAVEPOINT_CANCEL cancel请求参数为taskId与type。自定义参数的底层传递V1 提交时会调用generateVariables()构造variables字典其来源有两部分taskExecutionContext.getPrepareParamsMap()DolphinScheduler 的全局参数与流程级参数dinkyParameters.getLocalParams()节点上配置的局部参数其中会通过ParameterUtils.convertParameterPlaceholders把${param}占位符替换为实际值再合并进变量字典。最终这些变量以variables字段随submitTask请求发送给 Dinky 1.0实现跨系统参数透传。这也解释了文档中支持${param}方式获取全局或局部动态参数的约束——目前仅支持IN类型输入Dinky 侧暂不向 DolphinScheduler 回传OUT类型输出。适用前提与注意事项Dinky 版本要求从 DinkyTaskConstants.java 的API_VERSION_ERROR_TIPS可知插件要求Dinky 版本大于等于 0.6.5自定义参数特性则需 Dinky 1.0 及以上。网络连通性Worker 所在机器必须能访问 Dinky 服务的地址如http://localhost:8888且 Dinky Open API 已开放doGet/sendJsonStr仅在 HTTP 200 时才解析响应体网络不通或返回非 200 会直接导致任务失败。上线作业的前置条件勾选上线作业前请先在 Dinky 平台确认作业已发布且当前没有同名 Flink Job 实例在运行否则提交会失败。任务 ID 的唯一性Dinky 任务 ID必须是 Dinky 平台上真实存在的作业 ID配置错误时会在提交阶段报错并记录在 Worker 日志中。轮询开销默认每 3 秒轮询一次作业状态长时运行的 Flink 作业会产生持续的 HTTP 查询建议结合 DolphinScheduler 的失败重试与超时告警参数参见默认任务参数合理设定任务超时。小结Dinky 任务是 DolphinScheduler 连接实时计算生态的关键桥接插件之一它以DINKY任务类型的形式把 Dinky 平台的 FlinkSQL / Flink Jar / SQL 作业接入统一的工作流调度通过version → submitTask → getJobInstance → savepointTask四类 Open API 调用实现提交、跟踪与取消的完整生命周期管理并自动适配 Dinky 0.x/1.x 两套协议。理解其参数含义与源码实现可以帮助你在搭建离线调度 实时计算混合编排平台时准确配置节点、快速定位提交失败问题并借助自定义参数实现跨系统动态传参。【免费下载链接】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),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价