- 任务调度
- 大数据
- 后端
- 前端
【免费下载链接】dolphinscheduler
Apache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code
导读
本文围绕 Apache DolphinScheduler 中的 Sqoop 任务节点,系统讲解如何在低代码 DAG 中编排 Apache Sqoop 的 import / export 作业,实现 RDBMS 与 HDFS / Hive 之间的数据同步。你将掌握 Sqoop 节点的完整参数含义与取值规则、MySQL→Hive 导入的端到端配置示例,以及任务底层如何把表单参数动态拼装为真实sqoop命令行脚本(含源码佐证),从而在生产环境中可靠地使用该任务类型。
概览:Sqoop 任务节点是什么
Sqoop 任务类型用于执行 Sqoop 应用。DolphinScheduler 的 Worker 会调用本机环境中的sqoop命令来执行任务,因此它是一个"表单化包装 Sqoop 命令"的任务节点:用户在 DAG 画布上拖入 Sqoop 节点,以表单方式声明数据源、数据目标与同步参数,Worker 端最终将这些参数生成一段完整的sqoop脚本并以 YARN 任务的方式提交运行。
从源码结构看,该任务插件位于 dolphinscheduler-task-plugin/dolphinscheduler-task-sqoop,核心执行类 SqoopTask.java 继承自AbstractYarnTask(即它最终像 Spark、MapReduce 任务一样以 YARN 应用形态运行)。任务支持两种作业模式:
- TEMPLATE(模板模式):通过界面表单填写源/目标参数,由代码自动拼装
sqoop命令行,是本文重点; - CUSTOM(自定义模式):直接编写自定义
sqoopshell 脚本,适合需要完全控制命令行的场景。
创建 Sqoop 任务节点
在 DolphinScheduler Web UI 中创建 Sqoop 节点的操作路径如下:
- 进入
项目管理 -> 项目名称 -> 工作流定义,点击创建工作流按钮进入 DAG 编辑页面; - 从左侧工具栏拖拽 Sqoop 节点图标
到画布上(该图标即 Sqoop 节点,与官方文档中的图标一致)。
拖入画布后,双击节点即可在右侧配置面板中填写任务参数。
关于任务节点通用参数(如任务名称、运行环境、失败重试、告警分组、超时等),请参考 DolphinScheduler 任务参数附录 中的
Default Task Parameters一节,本文不再重复。
任务参数详解
Sqoop 节点在通用任务参数之上,还包含以下专属参数:
| 参数 | 说明 |
|---|---|
| Job Name | map-reduce 作业名称(对应-D mapred.job.name) |
| Direct | (1)import:将 RDBMS 中的单张表导入 HDFS 或 Hive;(2)export:将 HDFS 或 Hive 中的一组文件导出回 RDBMS |
| Hadoop Params | Sqoop 作业的自定义 Hadoop 参数(以-D key=value形式注入) |
| Sqoop Advanced Parameters | Sqoop 作业的高级参数(以key value形式原样追加到命令行) |
| Data Source - Type | 选择对应的数据源类型(MYSQL / HIVE / HDFS / ORACLE / HANA / SQLSERVER) |
| Data Source - Datasource | 选择对应的数据源(下拉选择在数据源中心预先注册的 DataSource) |
| Data Source - ModelType | (1)Form:从表同步数据,需填写Table和ColumnType;(2)SQL:同步 SQL 查询结果,需填写SQL Statement |
| Data Source - Table | 设置导入到 Hive 时使用的表名 |
| Data Source - ColumnType | (1)All Columns:导入所选表的全部字段;(2)Some Columns:导入所选表的指定字段,需填写Column |
| Data Source - Column | 填写字段名,多个字段用逗号分隔 |
| Data Source - SQL Statement | 填写 SQL 查询语句(对应--query) |
| Data Source - Map Column Hive | 覆盖指定列从 SQL 类型到 Hive 类型的映射(对应--map-column-hive) |
| Data Source - Map Column Java | 覆盖指定列从 SQL 类型到 Java 类型的映射(对应--map-column-java) |
| Data Target - Type | 选择对应的数据目标类型(MYSQL / HIVE / HDFS / ORACLE / HANA / SQLSERVER) |
| Data Target - Database | 填写 Hive 数据库名(对应--hive-database) |
| Data Target - Table | 填写 Hive 表名(对应--hive-table) |
| Data Target - CreateHiveTable | 是否在导入时于 Hive 中建表(对应--create-hive-table);若开启,当目标 Hive 表已存在时作业将失败 |
| Data Target - DropDelimiter | 导入 Hive 时是否丢弃字符串字段中的\n、\r和\01(对应--hive-drop-import-delims) |
| Data Target - OverWriteSrc | 是否覆盖 Hive 表中的已有数据(对应--hive-overwrite,同时追加--delete-target-dir) |
| Data Target - Hive Target Dir | 显式指定目标目录(对应--target-dir),不填则由 Sqoop 自动决定 |
| Data Target - ReplaceDelimiter | 导入 Hive 时用自定义字符串替换字段中的\n、\r、\01(对应--hive-delims-replacement) |
| Data Target - Hive partition Keys | 填写 Hive 分区键名,多个用逗号分隔(对应--hive-partition-key) |
| Data Target - Hive partition Values | 填写 Hive 分区值,多个用逗号分隔(对应--hive-partition-value) |
| Data Target - Target Dir | 填写 HDFS 目标目录(对应--target-dir) |
| Data Target - DeleteTargetDir | 若目标目录已存在则先删除(对应--delete-target-dir) |
| Data Target - CompressionCodec | 选择 Hadoop 压缩编解码器(对应--compression-codec) |
| Data Target - FileType | 选择存储类型(如 AVRO / PARQUET / TEXTFILE 等) |
| Data Target - FieldsTerminated | 设置字段分隔符(对应--fields-terminated-by) |
| Data Target - LinesTerminated | 设置行结束符(对应--lines-terminated-by) |
参数与底层命令的对应关系
上述表单参数并非直接透传,而是由插件内的生成器按角色拼装。从 SqoopConstants.java 可以看到所有常量与 Sqoop 原生命令的映射:
- 通用段:
sqoop、-D mapred.job.name、-D sqoop.export.records.per.statement、-m(并行度)、--split-by; - 数据库段:
--connect、--driver、--username、--password、--table、--columns、--query、--map-column-hive、--map-column-java; - HDFS 段:
--export-dir、--target-dir、--compression-codec; - Hive 段:
--hive-import、--hive-database、--hive-table、--create-hive-table、--hive-drop-import-delims、--hive-overwrite、--delete-target-dir、--hive-delims-replacement、--hive-partition-key、--hive-partition-value; - 更新模型:
--update-key、--update-mode。
数据源侧支持的类型在 SqoopParameters.java 的getSourceParameter/getTargetParameter中逐一映射到对应参数类,即源与目标均可选择 MYSQL、HIVE、HDFS、ORACLE、HANA、SQLSERVER 六种类型中的一种。
任务执行与脚本生成原理
执行链路
SqoopTask.java 的init()方法负责初始化:
- 将任务 JSON 参数反序列化为
SqoopParameters; - 调用
checkParameters()做参数合法性校验,失败则抛出TaskException; - 通过
generateExtendedContext()从资源参数中解析数据源连接信息(含解密后的连接参数); - 注册密码脱敏正则
SqoopConstants.SQOOP_PASSWORD_REGEX,确保--password "xxx"中的明文密码在日志中被打码显示。
随后getScript()调用 SqoopJobGenerator.java 生成脚本:
- TEMPLATE 模式:按
sourceType/targetType创建对应的 Source/Target 生成器,最终脚本 =CommonGenerator输出 +SourceGenerator输出 +TargetGenerator输出; - CUSTOM 模式:直接使用用户填写的
customShell(自动将\r\n统一为系统换行符)。
参数校验规则
SqoopParameters.checkParameters() 中规定了两种模式的必填约束:
- TEMPLATE:必须填写
modelType、jobName、sourceType、targetType、sourceParams、targetParams,且concurrency != 0、customShell为空; - CUSTOM:必须填写
customShell,且jobName为空。
注意concurrency != 0是模板模式的硬性校验项,因此在表单中并行度必须显式设置为非 0 值。
通用段脚本的生成细节
CommonGenerator.java 生成脚本的公共前缀,拼装顺序如下:
sqoop <modelType> -D sqoop.export.records.per.statement=1 -D mapred.job.name=<jobName>modelType即 import 或 export(Direct 参数);sqoop.export.records.per.statement固定为 1;- Hadoop Params中每一项以
-D prop=value追加; - Sqoop Advanced Parameters中每一项以
prop value原样追加(适合--connect、--driver这类键值对形式的原生参数); - 当
concurrency > 0时追加-m <concurrency>;当concurrency > 1时还会追加--split-by <splitBy>(split-by 由表单的拆分列决定)。
源端脚本生成细节(以 MySQL 为例)
MySQLSourceGenerator.java 负责拼装源端命令:
- 从任务上下文取出数据源连接参数,解密密码后生成:
--connect "jdbc:mysql://host:port/db" --username <user> --password "<decoded-password>"- 依据 ModelType 分支:
- Form:追加
--table <srcTable>;若 ColumnType 为 "Some Columns" 且已填列名,再追加--columns <col1,col2,...>; - SQL:追加
--query "<srcQuerySql> ...",并自动追加$CONDITIONS占位符——若 SQL 中已包含where关键字则追加AND \$CONDITIONS,否则追加WHERE \$CONDITIONS(这是 Sqoop 并发导入拆分数据的必要条件);
- Form:追加
- 若配置了Map Column Hive / Map Column Java,则按
prop=value,prop2=value2形式生成--map-column-hive .../--map-column-java ...。
其余源类型(HIVE / HDFS / ORACLE / HANA / SQLSERVER)均有对应的 sources 目录 下的独立生成器实现。
目标端脚本生成细节(以 Hive、HDFS 为例)
HiveTargetGenerator.java 生成 Hive 目标段命令:
- 始终以
--hive-import开头; - 填了库名和表名时追加
--hive-database <db> --hive-table <table>; CreateHiveTable=true追加--create-hive-table;DropDelimiter=true追加--hive-drop-import-delims;OverWriteSrc=true追加--hive-overwrite --delete-target-dir;- 填了
ReplaceDelimiter追加--hive-delims-replacement <value>; - 同时填了分区键与分区值时追加
--hive-partition-key <keys> --hive-partition-value <values>; - 填了
Hive Target Dir追加--target-dir <dir>。
HdfsTargetGenerator.java 则生成 HDFS 目标段:
- 依次追加
--target-dir <path>、--compression-codec <codec>、<fileType>、--delete-target-dir(可选); - 字段/行分隔符以单引号包裹:
--fields-terminated-by '<v>' --lines-terminated-by '<v>'; - 固定追加
--null-non-string 'NULL' --null-string 'NULL',将空值统一替换为 NULL 字符串。
数据源资源的绑定
模板模式下,任务会通过 getResources() 把源/目标数据源 ID 注册为ResourceType.DATASOURCE资源,执行时由 generateExtendedContext() 从资源中心解析出真实连接串并注入生成器。这意味着数据源密码等敏感信息不会出现在任务定义里,而是运行时统一获取、统一脱敏。
任务示例:从 MySQL 导入数据到 Hive
下面演示一个完整示例:将 MySQL 数据库test中的example表数据导入 Hive。示例数据如下图所示:
第一步:配置 Sqoop 运行环境
Sqoop 任务在 Worker 节点上以sqoop命令方式执行。若要在生产环境使用 Sqoop 任务类型,必须确保运行该任务的 Worker 机器上已正确安装 Sqoop,且sqoop命令在 PATH 中可用;同时 Worker 需具备访问目标 Hive / HDFS 集群的客户端环境(hive、hadoop 命令及对应依赖)。可先在 Worker 上手动执行sqoop version验证环境是否就绪。
第二步:配置 Sqoop 任务节点
按下图指引填写节点内容:
本示例的关键配置如下:
| 参数 | 值 |
|---|---|
| Job Name | sqoop_mysql_to_hive_test |
| Data Source - Type | MYSQL |
| Data Source - Datasource | MYSQL MyTestMySQL(MyTestMySQL 可替换为你自定义的数据源名称,需提前在数据源中心注册) |
| Data Source - ModelType | Form |
| Data Source - Table | example |
| Data Source - ColumnType | All Columns |
| Data Target - Type | HIVE |
| Data Target - Database | tmp |
| Data Target - Table | example |
| Data Target - CreateHiveTable | true |
| Data Target - DropDelimiter | false |
| Data Target - OverWriteSrc | true |
| Data Target - Hive Target Dir | (无需填写) |
| Data Target - ReplaceDelimiter | , |
| Data Target - Hive partition Keys | (无需填写) |
| Data Target - Hive partition Values | (无需填写) |
结合上文源码可知,该配置最终会生成类似如下的脚本(密码已被运行时注入与脱敏):
sqoop import -D sqoop.export.records.per.statement=1 -D mapred.job.name=sqoop_mysql_to_hive_test -m 1 \ --connect "jdbc:mysql://host:3306/test" --username <user> --password "******" \ --table example \ --hive-import --hive-database tmp --hive-table example \ --create-hive-table --hive-overwrite --delete-target-dir \ --hive-delims-replacement ,第三步:查看运行结果
保存并运行工作流后,可在任务实例/工作流实例页面查看执行日志与状态。任务成功时日志中会输出 MapReduce 作业的执行信息与导入行数统计:
若执行失败,优先排查三类问题:
- 环境类:Worker 上
sqoop命令缺失、Hive/HDFS 客户端配置错误或未认证(Kerberos)导致连接失败; - 数据源类:
Data Source - Datasource指向的数据源连接串、账号密码是否正确,Worker 是否能连通对应 RDBMS; - 目标表类:
CreateHiveTable=true时目标 Hive 表已存在会导致失败(正如参数说明所述),此时可改为 false 或先清理目标表。
常见使用建议
- 并行导入:需要提升导入性能时,将并行度(
-m)设为大于 1 的值并同时填写splitBy拆分列(源码会在concurrency > 1时自动追加--split-by); - SQL 模式必须带
$CONDITIONS:选择 SQL 模式时,插件会自动追加AND \$CONDITIONS或WHERE \$CONDITIONS,因此不要在该 SQL 末尾使用分号结尾,否则可能造成命令拼接异常; - 敏感信息保护:插件通过 SqoopConstants.SQOOP_PASSWORD_REGEX 对
--password "..."进行日志脱敏,排查问题时注意不要在自定义脚本中自行打印密码; - 自定义脚本场景:当表单无法覆盖某些高级 Sqoop 特性时,可切换到 CUSTOM 模式直接编写
sqoop脚本,此时 jobName 必须留空,脚本会原样执行。
小结
Sqoop 任务节点是 DolphinScheduler 数据同步工作流中的重要一环:它把繁琐的sqoop命令行封装为低代码表单,同时通过 SqoopTask.java →SqoopJobGenerator→ 各类 Source/Target 生成器的调用链,将表单参数逐段拼装为可在 YARN 上运行的 import / export 脚本。理解每一段命令的生成规则(通用段、源端段、目标端段)后,你就能准确预判任意表单组合产出的实际命令,从而快速定位导入导出失败的原因,并基于此设计出稳定、可维护的 MySQL→Hive(及 HDFS、Oracle、SQLServer、HANA 等)同步任务。
- 任务调度
- 大数据
- 后端
- 前端
【免费下载链接】dolphinscheduler
Apache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code
相关推荐
Apache DolphinScheduler Sqoop 任务节点:MySQL 到 Hive 数据导入实战与源码解析
Apache DolphinScheduler Sqoop 任务节点:MySQL 到 Hive 数据导入实战与源码解析 本文是 Apache DolphinSc
任务调度数据编排工作流自动化后端大数据Apache DolphinScheduler SQOOP 节点实战指南:任务参数详解与 MySQL 到 Hive 数据同步
Apache DolphinScheduler SQOOP 节点实战指南:任务参数详解与 MySQL 到 Hive 数据同步 本指南完整讲解 Apache Do
任务调度数据编排工作流自动化后端大数据Apache DolphinScheduler ChunJun 任务节点实战指南:从 JSON 配置到 Hive 数据同步
Apache DolphinScheduler ChunJun 任务节点实战指南:从 JSON 配置到 Hive 数据同步 ChunJun(原 FlinkX)是
任务调度大数据后端前端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考