ARTICLE DETAIL

资讯详情

深耕商务建站与企业官网运营的一线实战洞察。

DolphinScheduler DMS 节点实战指南:在数据编排工作流中创建、启动与重启 AWS DMS 迁移任务

DolphinScheduler DMS 节点实战指南:在数据编排工作流中创建、启动与重启 AWS DMS 迁移任务 DolphinScheduler DMS 节点实战指南在数据编排工作流中创建、启动与重启 AWS DMS 迁移任务【免费下载链接】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 内置DMSAWS Database Migration Service任务节点的完整使用与原理指南。AWS DMS 能在源数据库持续在线运行的前提下将数据快速、安全地迁移至 AWS 上的目标库是数据上云、异构库同步与 CDC变更数据捕获场景的常用服务。通过本节点的图形界面或 JSON 数据两种方式你可以把“创建迁移任务 启动迁移任务”“重启已有迁移任务”直接编排进 DolphinScheduler 的工作流 DAG并让调度平台自动跟踪任务状态直至迁移完成。读完本文你将掌握 DMS 节点全部参数含义、两种创建方式的实操步骤、AWS 环境配置方法以及组件状态跟踪与容错处理的底层实现逻辑。一、DMS 节点综述它能做什么AWS Database Migration Service 可帮助用户在数据源保持运行的情况下将数据迁移至 AWS从而最大程度减少依赖源数据库的应用程序停机时间该服务支持在广泛使用的开源与商业数据库之间迁移数据。DolphinScheduler 的 DMS 任务组件dolphinscheduler-task-plugin/dolphinscheduler-task-dms正是把这一能力封装成工作流中的一个节点使用户无需离开 DolphinScheduler 即可驱动 AWS DMS 迁移。组件主要包含两个功能创建并启动迁移任务在 AWS DMS 中新建一个 replication task并立即启动迁移流程重启已存在的迁移任务针对已有任务如迁移中断、需要 reload target 的场景仅执行启动/重启动作。组件的使用方式有两种通过界面创建在 DAG 编辑页中拖拽 DMS 节点在表单中逐项填写参数通过 JSON 数据创建将参数整体写入一段 JSON支持 DolphinScheduler 参数占位符替换适合模板化、批量化的任务定义。状态跟踪与“CDC 无结束时间”特例从源码看DolphinScheduler 启动 DMS 任务后并非“一放了之”而是持续跟踪任务状态直到任务完成才把工作流节点置为成功。其跟踪逻辑位于 DmsTask.java 的trackApplicationStatus()节点会先解析出 replication task 的 ARN再通过dmsHook.checkFinishedReplicationTask()轮询任务是否进入stopped状态且stopReason以FINISHED结尾详见 DmsHook.java。但存在一个明确例外不跟踪无结束时间的 CDC 任务。即当迁移类型为full-load-and-cdc或cdc且未配置cdcStopPosition参数时DolphinScheduler 在成功启动任务后直接将节点状态置为成功——因为这类 CDC 任务没有天然的结束时间点会持续运行等待增量数据。该判断由isStopTaskWhenCdc()实现DmsTask.java当实际任务的migrationType包含cdc且cdcStopPosition为空时返回 true节点随即成功退出。二、创建任务DAG 画布中的操作入口点击项目管理 → 项目名称 → 工作流定义点击“创建工作流”按钮进入 DAG 编辑页面从工具栏拖动 DMS 任务节点到画板中在右侧节点配置面板中设置节点名称、所属租户、运行环境等通用信息再按下文说明填写 DMS 专属参数。DMS 节点由DmsTaskChannelFactory注册组件名称为DMS见 DmsTaskChannelFactory.java前端在dolphinscheduler-ui/src/views/projects/task/components/node/tasks/index.ts中通过useDms完成表单渲染。表单会根据你选择的“是否重启任务 / 是否使用 JSON”动态显示/隐藏对应字段use-dms.ts这正好对应下文四种任务样例。三、任务样例四种典型配置3.1 创建并启动迁移任务通过界面适用于首次迁移填写源端点、目标端点、迁移实例 ARN 与迁移类型等参数DolphinScheduler 将依次调用 AWS DMS 的创建任务与启动任务接口。3.2 重启已存在的迁移任务通过界面适用于迁移中断或目标端需要重新加载数据的场景只需提供已有的replicationTaskArnDolphinScheduler 不会再创建新任务而是直接启动该任务默认以reload-target方式重启详见下文原理章节。3.3 创建并启动迁移任务通过 JSON 数据开启isJsonFormat后所有创建参数被整合进一段 JSON。JSON 中还可以使用file://前缀引用本地文件来承载大段的tableMappings与replicationTaskSettings内容。3.4 重启已存在的迁移任务通过 JSON 数据isJsonFormat与isRestartTask同时开启JSON 中仅需描述ReplicationTaskArn等重启所需字段。参数校验与“界面/JSON 两种方式”在服务端由DmsParameters.checkParameters()统一完成DmsParameters.javaJSON 方式要求jsonData非空重启方式要求replicationTaskArn非空创建方式则要求源/目标端点 ARN、迁移实例 ARN、迁移类型、任务标识符与表映射全部非空。四、参数详解4.1 DolphinScheduler 通用参数节点名称、运行环境、失败重试次数、超时告警、自定义参数localParams、资源resources等通用参数遵循所有任务节点共用的默认约定默认参数说明请参考 DolphinScheduler 任务参数附录 中“默认任务参数”一栏。其中自定义参数与资源同样出现在 DMS 节点的表单中use-dms.ts可用于实现参数化的迁移配置。4.2 DMS 组件独有参数参数名说明默认值isRestartTask是否重启已存在的迁移任务。开启后无需也不会创建新任务只需提供replicationTaskArnfalseisJsonFormat是否使用 JSON 格式的数据创建任务。开启后其余创建类参数被jsonData取代falsejsonDataJSON 格式的任务数据仅在isJsonFormat为 true 时生效无以上三个开关参数可在 DmsParameters.java 中看到其默认值与字段定义。4.3 创建并启动迁移任务时的参数参数名说明migrationType迁移类型可选值full-load全量、full-load-and-cdc全量 增量、cdc仅增量。前端下拉选项定义于 use-dms.tsreplicationTaskIdentifier迁移任务标识符即 AWS 侧的任务名称replicationInstanceArn迁移实例Replication Instance的 ARNsourceEndpointArn源端点的 ARNtargetEndpointArn目标端点的 ARNtableMappings表映射table mappingsJSON支持file://前缀引用本地文件除文档列出的上述必填项外从DmsParameters的字段定义DmsParameters.java还可以看到组件还透传了一批可选参数最终由DmsHook.createReplicationTask()原样写入 AWS 的CreateReplicationTaskRequestDmsHook.javareplicationTaskSettings复制任务设置任务级设置 JSON同样支持file://前缀cdcStartTimeCDC 起始时间cdcStartPositionCDC 起始位置cdcStopPositionCDC 停止位置配置后任务可自动结束否则 DolphinScheduler 不再跟踪见上文 CDC 特例tagsAWS 资源标签列表taskData/resourceIdentifierAWS DMS 补充任务数据与资源标识符。4.4 重启已存在的迁移任务时的参数参数名说明replicationTaskArn待重启的迁移任务的 ARN4.5 JSON 数据示例与字段命名规则JSON 数据最终被反序列化为DmsParameters反序列化采用UpperCamelCaseStrategy命名策略DmsTask.java因此 JSON 中的字段名必须使用大驼峰形式与界面表单中的小驼峰字段一一对应。参考 DmsTaskTest.java 中的测试样例一个完整的“创建并启动”JSON 如下{ ReplicationTaskIdentifier: task6, SourceEndpointArn: arn:aws:dms:ap-southeast-1:511640773671:endpoint:Z7SUEAL273SCT7OCPYNF5YNDHJDDFRATGNQISOQ, TargetEndpointArn: arn:aws:dms:ap-southeast-1:511640773671:endpoint:aws-mysql57-target, ReplicationInstanceArn: arn:aws:dms:ap-southeast-1:511640773671:rep:dms2c2g, MigrationType: full-load, TableMappings: file://table-mapping.json, ReplicationTaskSettings: file://ReplicationTaskSettings.json, Tags: [ { Key: key1, Value: value1 } ] }几点说明TableMappings与ReplicationTaskSettings中的file://前缀由 DmsHook.replaceFileParameters() 解析即从本地文件读取内容替换为实际 JSON 字符串适合承载超长映射配置JSON 中支持 DolphinScheduler 参数占位符convertJsonParameters()会先通过ParameterUtils.convertParameterPlaceholders()完成占位符替换再执行反序列化DmsTask.java因此可以将日期、表名等动态值注入迁移配置反序列化时FAIL_ON_UNKNOWN_PROPERTIES被关闭JSON 中额外字段不会导致失败。五、环境配置连接 AWS 所需的 yaml 设置使用 DMS 节点前需要为 AWS 配置访问凭证。修改项目中的aws.yaml模板位于 dolphinscheduler-authentication/dolphinscheduler-aws-authentication/src/main/resources/aws.yaml生产环境按部署方式挂载到各服务的配置目录中aws.dms一段dms: # The AWS credentials provider type. support: AWSStaticCredentialsProvider, InstanceProfileCredentialsProvider # AWSStaticCredentialsProvider: use the access key and secret key to authenticate # InstanceProfileCredentialsProvider: use the IAM role to authenticate credentials.provider.type: AWSStaticCredentialsProvider access.key.id: access.key.id access.key.secret: access.key.secret region: region endpoint: endpoint各配置项含义credentials.provider.type凭证提供者类型。AWSStaticCredentialsProvider使用 Access Key / Secret Key 静态认证InstanceProfileCredentialsProvider使用 IAM 角色认证适用于运行在 AWS 上的实例。该键在 AwsConfigurationKeys.java 中定义非法值会在AWSCredentialsProviderFactor中被拒绝access.key.id / access.key.secret静态认证时的访问密钥regionDMS 服务所在区域endpoint服务端点可配置为 AWS 官方端点或自建兼容端点如本地模拟环境。运行时DmsHook.createClient()通过PropertyUtils.getByPrefix(aws.dms., )读取上述前缀配置并交由AWSDatabaseMigrationServiceClientFactory构建 AWS DMS 客户端DmsHook.java。注意配置前缀是aws.dms.而默认 yaml 中对应的就是aws.dms.credentials.provider.type、aws.dms.access.key.id等键。六、运行原理从提交到跟踪的完整链路6.1 任务提交submitApplicationDmsTask.submitApplication()DmsTask.java的执行顺序为创建任务若isRestartTask为 false调用dmsHook.createReplicationTask()创建 replication task并轮询等待其进入ready状态若为 true 则跳过创建步骤启动任务调用dmsHook.startReplicationTask()轮询等待任务进入running状态。启动类型由initDmsHook()决定重启任务默认reload-target重载目标端数据后启动新建任务默认start-replicationDmsTask.java也支持通过startReplicationTaskType参数覆盖失败清理若启动失败且不是重启场景调用dmsHook.deleteReplicationTask()删除刚创建的任务避免残留ResourceNotFoundException时视为已删除成功记录 ARN成功后把 replication task 的 ARN 保存为appIdApplicationIds结构供后续状态跟踪与日志查看使用。轮询采用awaitReplicationTaskStatus每CHECK_INTERVAL 1000ms 查询一次并在任务处于running/stopped时输出fullLoadProgressPercent全量加载进度百分比DmsHook.java。6.2 启动容错端点连接测试与自动重试启动阶段如果 AWS 抛出InvalidResourceStateException组件会检查错误信息是否包含Test connection意味着迁移实例暂时无法连通源/目标端点若错误与连接无关直接判定失败若为连接类错误则先调用testConnectionEndpoint()依次测试迁移实例与源、目标端点的连通性轮询DescribeConnections直至状态为successful测试通过后再次尝试启动任务DmsTask.java。这一逻辑在 DmsTaskTest.testStartReplicationTaskRestartTestConnection 中有对应用例覆盖。6.3 状态跟踪与取消工作流运行期间trackApplicationStatus()轮询任务状态任务进入stopped后检查stopReason以FINISHED结尾视为成功否则抛出TaskException标记节点失败当工作流被取消如人工 kill 或超时终止时cancelApplication()调用dmsHook.stopReplicationTask()停止正在运行的 AWS 迁移任务DmsTask.java避免迁移任务在 AWS 侧继续空转。6.4 状态机常量一览DmsHook.java 集中定义了组件使用的 AWS DMS 状态常量便于理解日志与跟踪行为常量值用途STATUS.DELETEdelete删除任务时的期望状态STATUS.READYready创建任务后等待的就绪状态STATUS.RUNNINGrunning启动任务后等待的运行状态STATUS.STOPPEDstopped任务结束后的状态STATUS.SUCCESSFUL/STATUS.TESTINGsuccessful/testing端点连通性测试的状态STATUS.FINISH_END_TOKENFINISHEDstopReason的成功结束标记七、FAQ 与注意事项迁移任务一直不结束节点也一直 running确认迁移类型是否为full-load-and-cdc/cdc且未设置cdcStopPosition。若是属于预期的“不跟踪 CDC”行为节点会在启动成功后立即置为成功可结合isStopTaskWhenCdc()的日志确认若需要任务有结束点请显式配置cdcStopPosition。JSON 方式创建任务报参数错误检查 JSON 字段是否使用大驼峰命名如ReplicationTaskIdentifier且isJsonFormat为 true 时jsonData必须非空。启动时报 “Test connection” 错误组件会自动重测源/目标端点连通性并重试若持续失败请先在 AWS 控制台确认 Replication Instance 与两端点的网络、安全组与连接配置。JSON 中如何引用大段表映射使用file://前缀指向本地文件如file://table-mapping.json组件会读取文件内容作为实际参数请确保该文件在 DolphinScheduler Worker 节点上可访问。误创建的任务会不会残留启动失败时组件会自动删除刚创建的 replication task若删除返回ResourceNotFoundException亦视为已删除。以上行为均可对照源码验证节点参数模型见 DmsParameters.javaAWS 交互封装见 DmsHook.java任务生命周期编排见 DmsTask.java对应测试覆盖在 DmsTaskTest.java 与 DmsHookTest.java 中。/DSMLparameter /DSMLinvoke /DSMLtool_calls【免费下载链接】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),仅供参考
返回列表
PREV
查看更多资讯
NEXT
返回资讯列表