Apache DolphinScheduler Sqoop 任务节点实战指南:从 MySQL 到 Hive 的数据导入
2026/9/23 9:46:57 网站建设 项目流程
  • 任务调度
  • 大数据
  • 后端
  • 前端

【免费下载链接】dolphinscheduler

Apache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code

项目地址:https://gitcode.com/gh_mirrors/do/dolphinscheduler
点击查看免费下载

导读

本文围绕 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 节点的操作路径如下:

  1. 进入项目管理 -> 项目名称 -> 工作流定义,点击创建工作流按钮进入 DAG 编辑页面;
  2. 从左侧工具栏拖拽 Sqoop 节点图标到画布上(该图标即 Sqoop 节点,与官方文档中的图标一致)。

拖入画布后,双击节点即可在右侧配置面板中填写任务参数。

关于任务节点通用参数(如任务名称、运行环境、失败重试、告警分组、超时等),请参考 DolphinScheduler 任务参数附录 中的Default Task Parameters一节,本文不再重复。

任务参数详解

Sqoop 节点在通用任务参数之上,还包含以下专属参数:

参数说明
Job Namemap-reduce 作业名称(对应-D mapred.job.name
Direct(1)import:将 RDBMS 中的单张表导入 HDFS 或 Hive;(2)export:将 HDFS 或 Hive 中的一组文件导出回 RDBMS
Hadoop ParamsSqoop 作业的自定义 Hadoop 参数(以-D key=value形式注入)
Sqoop Advanced ParametersSqoop 作业的高级参数(以key value形式原样追加到命令行)
Data Source - Type选择对应的数据源类型(MYSQL / HIVE / HDFS / ORACLE / HANA / SQLSERVER)
Data Source - Datasource选择对应的数据源(下拉选择在数据源中心预先注册的 DataSource)
Data Source - ModelType(1)Form:从表同步数据,需填写TableColumnType;(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()方法负责初始化:

  1. 将任务 JSON 参数反序列化为SqoopParameters
  2. 调用checkParameters()做参数合法性校验,失败则抛出TaskException
  3. 通过generateExtendedContext()从资源参数中解析数据源连接信息(含解密后的连接参数);
  4. 注册密码脱敏正则SqoopConstants.SQOOP_PASSWORD_REGEX,确保--password "xxx"中的明文密码在日志中被打码显示。

随后getScript()调用 SqoopJobGenerator.java 生成脚本:

  • TEMPLATE 模式:按sourceType/targetType创建对应的 Source/Target 生成器,最终脚本 =CommonGenerator输出 +SourceGenerator输出 +TargetGenerator输出;
  • CUSTOM 模式:直接使用用户填写的customShell(自动将\r\n统一为系统换行符)。

参数校验规则

SqoopParameters.checkParameters() 中规定了两种模式的必填约束:

  • TEMPLATE:必须填写modelTypejobNamesourceTypetargetTypesourceParamstargetParams,且concurrency != 0customShell为空;
  • 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 负责拼装源端命令:

  1. 从任务上下文取出数据源连接参数,解密密码后生成:
--connect "jdbc:mysql://host:port/db" --username <user> --password "<decoded-password>"
  1. 依据 ModelType 分支:
    • Form:追加--table <srcTable>;若 ColumnType 为 "Some Columns" 且已填列名,再追加--columns <col1,col2,...>
    • SQL:追加--query "<srcQuerySql> ...",并自动追加$CONDITIONS占位符——若 SQL 中已包含where关键字则追加AND \$CONDITIONS,否则追加WHERE \$CONDITIONS(这是 Sqoop 并发导入拆分数据的必要条件);
  2. 若配置了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 Namesqoop_mysql_to_hive_test
Data Source - TypeMYSQL
Data Source - DatasourceMYSQL MyTestMySQL(MyTestMySQL 可替换为你自定义的数据源名称,需提前在数据源中心注册)
Data Source - ModelTypeForm
Data Source - Tableexample
Data Source - ColumnTypeAll Columns
Data Target - TypeHIVE
Data Target - Databasetmp
Data Target - Tableexample
Data Target - CreateHiveTabletrue
Data Target - DropDelimiterfalse
Data Target - OverWriteSrctrue
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 作业的执行信息与导入行数统计:

若执行失败,优先排查三类问题:

  1. 环境类:Worker 上sqoop命令缺失、Hive/HDFS 客户端配置错误或未认证(Kerberos)导致连接失败;
  2. 数据源类Data Source - Datasource指向的数据源连接串、账号密码是否正确,Worker 是否能连通对应 RDBMS;
  3. 目标表类CreateHiveTable=true时目标 Hive 表已存在会导致失败(正如参数说明所述),此时可改为 false 或先清理目标表。

常见使用建议

  • 并行导入:需要提升导入性能时,将并行度(-m)设为大于 1 的值并同时填写splitBy拆分列(源码会在concurrency > 1时自动追加--split-by);
  • SQL 模式必须带$CONDITIONS:选择 SQL 模式时,插件会自动追加AND \$CONDITIONSWHERE \$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

项目地址:https://gitcode.com/gh_mirrors/do/dolphinscheduler
点击查看免费下载

相关推荐

上一篇:SpatialThinker-30B-i1-GGUF模型架构详解:专家混合与多模态融合技术
下一篇:【亲测免费】 推荐开源项目:《学习 Go 语言》第二版

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询