5.6 工作流调度工具 Oozie 第五章:Hadoop 生态系统工具与应用 5.6 工作流调度工具 Oozie 在庞大的 Hadoop 生态系统中,数据处理流程往往错综复杂,涉及多个步骤和技术组件的协同工作。从数据抽取、转换、加载(ETL),到数据分析、机器学习模型训练,再到结果可视化,一个完整的数据管道可能包含 MapReduce、Pig、Hive、Spark 等多种任务类型。如何有效地组织、调度和监控这些复杂的任务流程,确保数据处理的可靠性和效率,成为了 Hadoop 应用开发中的关键挑战。为了解决这一问题,Apache Oozie 应运而生,成为了 Hadoop 生态系统中首屈一指的工作流调度工具。 5.6.1 Oozie 概述 什么是 Oozie?
在庞大的 Hadoop 生态系统中,数据处理流程往往错综复杂,涉及多个步骤和技术组件的协同工作。从数据抽取、转换、加载(ETL),到数据分析、机器学习模型训练,再到结果可视化,一个完整的数据管道可能包含 MapReduce、Pig、Hive、Spark 等多种任务类型。如何有效地组织、调度和监控这些复杂的任务流程,确保数据处理的可靠性和效率,成为了 Hadoop 应用开发中的关键挑战。为了解决这一问题,Apache Oozie 应运而生,成为了 Hadoop 生态系统中首屈一指的工作流调度工具。
什么是 Oozie?
Apache Oozie 是一个开源的工作流调度系统,用于管理 Hadoop 集群上的工作流。它可以将多个 Hadoop 任务(如 MapReduce、Pig、Hive、Spark 等)以及 Shell 脚本、Java 程序等,按照预定义的顺序和依赖关系,组织成一个有向无环图(DAG)的工作流。Oozie 负责协调这些任务的执行,处理任务之间的依赖关系,监控任务的运行状态,并在任务失败时进行重试或告警,从而实现自动化、可靠的工作流管理。
为什么需要 Oozie?
在 Hadoop 环境中,手动管理复杂的任务流程是低效且容易出错的。例如,一个典型的 ETL 流程可能包含以下步骤:
数据抽取: 从外部数据源(如关系数据库、日志文件等)抽取数据到 HDFS。
数据清洗和转换: 使用 Pig 或 Spark 对数据进行清洗、转换和预处理。
数据加载: 将清洗后的数据加载到 Hive 数据仓库或 HBase 数据库。
数据分析: 使用 Hive 或 Spark SQL 进行数据分析和报表生成。
如果这些步骤需要手动依次执行,不仅耗时耗力,而且难以维护和监控。一旦某个步骤失败,就需要人工介入排查问题并重新启动任务。而 Oozie 的出现,正是为了解决这些痛点:
自动化调度: Oozie 可以根据预定义的调度策略,自动触发工作流的执行,无需人工干预。
任务编排: Oozie 允许用户将多个任务按照依赖关系组织成工作流,实现复杂的任务编排。
可靠性保证: Oozie 具有完善的错误处理机制,可以在任务失败时进行重试、告警或执行补偿操作,确保工作流的可靠运行。
集中监控: Oozie 提供了 Web UI 和命令行工具,方便用户监控工作流的运行状态、日志信息和执行历史。
灵活性和扩展性: Oozie 支持多种任务类型,并可以通过扩展 Action 节点来支持自定义任务,具有良好的灵活性和扩展性。
Oozie 的主要特点:
基于 XML 的工作流定义: Oozie 工作流和协调器都是使用 XML 文件定义的,易于理解和维护。
支持多种任务类型: Oozie 可以调度 MapReduce、Pig、Hive、Spark、Shell 脚本、Java 程序等多种类型的任务。
时间驱动和数据驱动的调度: Oozie 支持基于时间的调度(使用 Coordinator)和基于数据可用性的调度。
强大的错误处理机制: Oozie 支持任务重试、Kill 节点、错误捕获和告警等多种错误处理机制。
Web UI 和命令行接口: Oozie 提供了友好的 Web UI 和强大的命令行接口,方便用户管理和监控工作流。
与 Hadoop 生态系统深度集成: Oozie 与 Hadoop 生态系统中的其他组件(如 HDFS、YARN、Hive、Pig 等)深度集成,可以无缝地调度各种 Hadoop 任务。
Oozie 的架构主要由以下几个核心组件构成:
Oozie Server: Oozie 服务器是 Oozie 的核心组件,负责接收、解析、调度和监控工作流。Oozie Server 运行在 Hadoop 集群中的一台或多台服务器上,通常部署在单独的节点上以保证其稳定性。它主要包含以下功能模块:
Workflow Engine: 工作流引擎负责解析工作流定义文件(workflow.xml),构建工作流的有向无环图,并根据工作流的定义,调度和执行工作流中的各个 Action 节点。
Coordinator Engine: 协调器引擎负责解析协调器定义文件(coordinator.xml),根据时间或数据依赖关系,生成和调度多个工作流实例。
JMS Server: Oozie 使用 JMS(Java Message Service)作为消息队列,用于组件之间的异步通信,例如 Oozie Client 与 Oozie Server 之间的通信,以及 Action 节点与 Oozie Server 之间的状态更新。
Database: Oozie 使用数据库(通常是 Derby、MySQL 或 PostgreSQL)来持久化工作流定义、工作流实例状态、日志信息等元数据。
Oozie Client: Oozie 客户端是用户与 Oozie Server 交互的工具,提供了命令行接口(CLI)和 Web UI 两种方式。
Oozie CLI: Oozie 命令行客户端允许用户通过命令行提交、监控和管理工作流和协调器。用户可以使用 oozie 命令来执行各种操作,如提交工作流、查看工作流状态、杀死工作流等。
Oozie Web UI: Oozie Web UI 提供了一个图形化的界面,方便用户查看工作流的运行状态、日志信息、执行历史等。Web UI 通常部署在 Oozie Server 所在的节点上,可以通过浏览器访问。
Workflow Definition (workflow.xml): 工作流定义文件是描述工作流逻辑的核心文件,使用 XML 格式编写。它定义了工作流包含的 Action 节点、Control 节点以及节点之间的连接关系。Workflow Definition 文件描述了任务的执行顺序和依赖关系,是 Oozie 调度工作流的依据。
Coordinator Definition (coordinator.xml): 协调器定义文件用于定义基于时间或数据触发的工作流调度策略,也使用 XML 格式编写。Coordinator Definition 文件描述了工作流的调度频率、时间范围、输入数据集和输出数据集等信息,Oozie Coordinator Engine 根据这些信息生成和调度工作流实例。
Action Nodes: Action 节点是工作流中的基本执行单元,代表一个具体的任务,例如执行 MapReduce 作业、Pig 脚本、Hive 查询、Spark 应用、Shell 脚本或 Java 程序等。Oozie 提供了多种内置的 Action 类型,用户也可以自定义 Action 类型。
Control Nodes: Control 节点用于控制工作流的执行流程,包括 start、end、decision、fork、join 和 kill 等节点。Control 节点不执行具体的任务,而是用于定义工作流的逻辑结构。
Oozie 架构示意图:
工作流定义文件 workflow.xml 是 Oozie 的核心,它使用 XML 语法描述了工作流的逻辑结构和任务流程。一个 workflow.xml 文件主要包含以下元素:
<workflow-app>: 根元素,定义工作流应用程序。
name: 工作流应用程序的名称,在 Oozie 中唯一标识一个工作流。
xmlns: XML 命名空间,通常为 uri:oozie:workflow:0.5。
<start>: 定义工作流的起始节点。
to: 指定工作流启动后要执行的第一个节点名称。<end>: 定义工作流的结束节点,表示工作流成功完成。
name: 结束节点的名称,通常命名为 end。<action>: 定义一个 Action 节点,代表一个具体的任务。
name: Action 节点的名称,在工作流中唯一标识一个 Action。
cred: 可选属性,指定用于执行 Action 的凭据(如 Kerberos 票据)。
Action 类型元素(如 <map-reduce>, <pig>, <hive>, <spark>, <shell>, <java> 等): 定义 Action 的具体类型和配置信息。
<ok to="节点名称"/>: 指定 Action 成功完成后要执行的下一个节点。
<error to="节点名称"/>: 指定 Action 执行失败后要执行的节点,通常指向 kill 节点。
<decision>: 定义一个决策节点,根据条件判断选择不同的执行路径。
name: 决策节点的名称。
<switch>: 包含多个 <case> 和一个 <default> 子元素。
<case condition="EL 表达式" to="节点名称"/>: 定义一个条件分支,如果条件满足(EL 表达式求值为 true),则跳转到指定的节点。
<default to="节点名称"/>: 定义默认分支,如果所有 <case> 条件都不满足,则跳转到指定的节点。
<fork>: 定义一个分支节点,用于并行执行多个 Action 节点。
name: 分支节点的名称。
<path start="节点名称"/>: 定义一个分支路径,指定该分支的起始节点。
<join>: 定义一个汇合节点,用于等待所有并行分支执行完成后再继续执行后续节点。
name: 汇合节点的名称。
to: 指定所有分支汇合后要执行的下一个节点名称。
fork: 指定要汇合的分支节点的名称,必须与 fork 节点的名称一致。
<kill>: 定义一个 Kill 节点,用于终止工作流的执行,通常在发生错误时跳转到 Kill 节点。
name: Kill 节点的名称,通常命名为 kill。
<message>错误信息</message>: 可选元素,指定 Kill 节点的错误信息,可以在 Oozie Web UI 中查看。
Action 类型元素:
Oozie 支持多种 Action 类型,常用的 Action 类型包括:
<map-reduce>: 用于执行 MapReduce 作业。
<job-tracker>: JobTracker 地址。
<name-node>: NameNode 地址。
<configuration>: MapReduce 作业的配置信息,例如输入路径、输出路径、Mapper 类、Reducer 类等。
<pig>: 用于执行 Pig 脚本。
<job-tracker>: JobTracker 地址。
<name-node>: NameNode 地址。
<script>: Pig 脚本的路径。
<param>: 传递给 Pig 脚本的参数。
<hive>: 用于执行 Hive 查询。
<job-tracker>: JobTracker 地址。
<name-node>: NameNode 地址。
<script>: Hive 查询脚本的路径。
<param>: 传递给 Hive 查询脚本的参数。
<spark>: 用于执行 Spark 应用。
<job-tracker>: JobTracker 地址(YARN 模式下可以省略)。
<name-node>: NameNode 地址。
<master>: Spark Master 地址(YARN 模式下为 yarn)。
<mode>: Spark 部署模式(cluster 或 client,YARN 模式下通常为 cluster)。
<name>: Spark 应用名称。
<class>: Spark 应用的主类。
<jar>: Spark 应用的 JAR 文件路径。
<arg>: 传递给 Spark 应用的参数。
<shell>: 用于执行 Shell 脚本。
<exec>: 要执行的 Shell 命令或脚本路径。
<argument>: 传递给 Shell 命令或脚本的参数。
<env-var>: 设置环境变量。
<capture-output>: 是否捕获 Shell 命令的输出到 Oozie 日志中。
<java>: 用于执行 Java 程序。
<job-tracker>: JobTracker 地址。
<name-node>: NameNode 地址。
<main-class>: Java 程序的主类。
<jar>: Java 程序的 JAR 文件路径。
<arg>: 传递给 Java 程序的参数。
工作流定义示例 (workflow.xml):
以下示例展示了一个简单的工作流,包含两个 Action 节点:一个 Shell 脚本节点和一个 Hive 查询节点。
<workflow-app xmlns="uri:oozie:workflow:0.5" name="shell-hive-workflow"> <start to="shell-action"/> <action name="shell-action"> <shell xmlns="uri:oozie:shell-action:0.2"> <exec>echo</exec> <argument>Hello from Shell Action!</argument> </shell> <ok to="hive-action"/> <error to="kill"/> </action> <action name="hive-action"> <hive xmlns="uri:oozie:hive-action:0.2"> <job-tracker>${jobTracker}</job-tracker> <name-node>${nameNode}</name-node> <configuration> <property> <name>mapred.job.queue.name</name> <value>${queueName}</value> </property> </configuration> <script>scripts/hive_script.hql</script> <param>inputPath=/user/${wf:user()}/input</param> <param>outputPath=/user/${wf:user()}/output</param> </hive> <ok to="end"/> <error to="kill"/> </action> <kill name="kill"> <message>Workflow failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message> </kill> <end name="end"/> </workflow-app>
代码详解:
<workflow-app> 根元素定义了工作流应用程序,名称为 shell-hive-workflow。
<start to="shell-action"/> 定义了工作流的起始节点为 shell-action。
<action name="shell-action"> 定义了一个名为 shell-action 的 Action 节点,类型为 <shell>,执行 echo "Hello from Shell Action!" 命令。
<ok to="hive-action"/> 表示 shell-action 成功完成后跳转到 hive-action 节点。
<error to="kill"/> 表示 shell-action 失败后跳转到 kill 节点。
<action name="hive-action"> 定义了一个名为 hive-action 的 Action 节点,类型为 <hive>,执行 Hive 脚本 scripts/hive_script.hql。
<job-tracker>${jobTracker}</job-tracker> 和 <name-node>${nameNode}</name-node> 使用了 EL 表达式 ${} 引用外部属性文件中的 jobTracker 和 nameNode 属性值。
<configuration> 定义了 Hive 作业的配置信息,例如队列名称。
<script>scripts/hive_script.hql</script> 指定了 Hive 脚本的路径。
<param> 定义了传递给 Hive 脚本的参数 inputPath 和 outputPath,也使用了 EL 表达式 ${wf:user()} 获取当前 Oozie 用户。
<ok to="end"/> 表示 hive-action 成功完成后跳转到 end 节点。
<error to="kill"/> 表示 hive-action 失败后跳转到 kill 节点。
<kill name="kill"> 定义了一个名为 kill 的 Kill 节点,用于终止工作流,并显示错误信息。
<end name="end"/> 定义了工作流的结束节点 end,表示工作流成功完成。
mermaid 图示工作流流程:
协调器定义文件 coordinator.xml 用于定义基于时间或数据触发的工作流调度策略。一个 coordinator.xml 文件主要包含以下元素:
<coordinator-app>: 根元素,定义协调器应用程序。
name: 协调器应用程序的名称,在 Oozie 中唯一标识一个协调器。
frequency: 定义工作流实例的调度频率,可以使用 Cron 表达式或时间间隔表达式(如 ${hours(1)} 表示每小时)。
start: 定义协调器开始调度工作流实例的时间,格式为 yyyy-MM-ddTHH:mmZ(UTC 时间)。
end: 定义协调器结束调度工作流实例的时间,格式同 start。
timezone: 定义协调器使用的时间zone,默认为 UTC。
xmlns: XML 命名空间,通常为 uri:oozie:coordinator:0.4。
<controls>: 定义协调器的控制参数。
<timeout>: 定义工作流实例的超时时间,单位为分钟。
<concurrency>: 定义协调器允许并行运行的最大工作流实例数量。
<execution>: 定义工作流实例的执行顺序,可以是 FIFO(先进先出)或 LIFO(后进先出)。
<datasets>: 定义协调器依赖的数据集。
<dataset name="数据集名称" frequency="数据集更新频率" initial-instance="数据集初始实例时间" timezone="时区">
<uri-prefix>: 数据集 URI 的前缀。
<uri-suffix>: 数据集 URI 的后缀。
<done-flag>: 数据集完成标志文件,可选。
<input-events>: 定义协调器的输入事件,用于描述工作流实例的输入数据依赖。
<data-in name="事件名称" dataset="数据集名称">
<instance>${coord:current(n)}</instance>: 指定输入数据集的实例,可以使用 EL 表达式 ${coord:current(n)} 获取数据集的第 n 个实例(n=0 表示当前实例,n=-1 表示上一个实例,以此类推)。<output-events>: 定义协调器的输出事件,用于描述工作流实例的输出数据。
<data-out name="事件名称" dataset="数据集名称">
<instance>${coord:current(0)}</instance>: 指定输出数据集的实例,通常使用 ${coord:current(0)} 表示当前实例。<action>: 定义协调器调度的 Action,通常是一个工作流 Action。
<workflow>: 定义要调度的工作流。
<app-path>: 工作流应用程序的 HDFS 路径。
<configuration>: 传递给工作流的配置参数,可以使用 <property> 元素定义。
协调器定义示例 (coordinator.xml):
以下示例展示了一个基于时间的协调器,每小时调度一次 shell-hive-workflow 工作流。
<coordinator-app name="hourly-shell-hive-coordinator" frequency="${coord:hours(1)}" start="${startTime}" end="${endTime}" timezone="UTC" xmlns="uri:oozie:coordinator:0.4"> <controls> <timeout>60</timeout> <concurrency>1</concurrency> <execution>FIFO</execution> </controls> <datasets> <dataset name="input-data" frequency="${coord:hours(1)}" initial-instance="${startTime}" timezone="UTC"> <uri-prefix>/user/${wf:user()}/input-data/</uri-prefix> <uri-suffix>/*</uri-suffix> </dataset> <dataset name="output-data" frequency="${coord:hours(1)}" initial-instance="${startTime}" timezone="UTC"> <uri-prefix>/user/${wf:user()}/output-data/</uri-prefix> <uri-suffix>/*</uri-suffix> </dataset> </datasets> <input-events> <data-in name="input" dataset="input-data"> <instance>${coord:current(-1)}</instance> </data-in> </input-events> <output-events> <data-out name="output" dataset="output-data"> <instance>${coord:current(0)}</instance> </data-out> </output-events> <action> <workflow> <app-path>${workflowAppUri}</app-path> <configuration> <property> <name>queueName</name> <value>${queueName}</value> </property> <property> <name>inputPath</name> <value>${coord:dataIn('input')}</value> </property> <property> <name>outputPath</name> <value>${coord:dataOut('output')}</value> </property> </configuration> </workflow> </action> </coordinator-app>
代码详解:
<coordinator-app> 根元素定义了协调器应用程序,名称为 hourly-shell-hive-coordinator。
frequency="${coord:hours(1)}" 定义调度频率为每小时一次。
start="${startTime}" 和 end="${endTime}" 定义了协调器的开始和结束时间,使用外部属性 ${startTime} 和 ${endTime}。
timezone="UTC" 指定时区为 UTC。
<controls> 定义了协调器的控制参数,例如超时时间、并发数和执行顺序。
<datasets> 定义了两个数据集 input-data 和 output-data,都设置为每小时更新一次。
<input-events> 定义了一个输入事件 input,依赖于 input-data 数据集的上一个实例 ${coord:current(-1)}。
<output-events> 定义了一个输出事件 output,依赖于 output-data 数据集的当前实例 ${coord:current(0)}。
<action> 定义了要调度的 Action,类型为 <workflow>,即调度一个工作流。
<app-path>${workflowAppUri}</app-path> 指定了工作流应用程序的 HDFS 路径,使用外部属性 ${workflowAppUri}。
<configuration> 定义了传递给工作流的配置参数。
queueName 和 workflowAppUri 来自外部属性。
inputPath 使用 ${coord:dataIn('input')} EL 表达式获取输入事件 input 对应的数据集 URI。
outputPath 使用 ${coord:dataOut('output')} EL 表达式获取输出事件 output 对应的数据集 URI。
mermaid 图示协调器调度流程:
1. 准备工作流应用程序:
首先,需要将工作流定义文件 workflow.xml、Hive 脚本 scripts/hive_script.hql 以及其他必要的资源文件(如 JAR 包、配置文件等)上传到 HDFS 上的一个目录中,例如 /user/${USER}/oozie/workflows/shell-hive-workflow/。
2. 创建 job.properties 文件:
创建一个 job.properties 文件,用于配置工作流的参数,例如 JobTracker 地址、NameNode 地址、队列名称、工作流应用程序路径等。
oozie.wf.application.path=${nameNode}/user/${USER}/oozie/workflows/shell-hive-workflow/workflow.xml oozie.coord.application.path=${nameNode}/user/${USER}/oozie/coordinators/hourly-shell-hive-coordinator/coordinator.xml oozie.use.system.libpath=true nameNode=hdfs://namenode:8020 jobTracker=yarn:resourcemanager:8032 queueName=default startTime=2024-01-01T00:00Z endTime=2024-01-02T00:00Z workflowAppUri=${nameNode}/user/${USER}/oozie/workflows/shell-hive-workflow/workflow.xml
3. 提交工作流或协调器:
使用 Oozie 命令行客户端 oozie 命令提交工作流或协调器。
oozie job -oozie http://oozie-server:11000/oozie -config job.properties -run
oozie job -oozie http://oozie-server:11000/oozie -config job.properties -run
4. 监控工作流或协调器状态:
可以使用 Oozie 命令行客户端或 Web UI 监控工作流或协调器的状态。
oozie job -oozie http://oozie-server:11000/oozie -info <job_id>
将 <job_id> 替换为实际的工作流 Job ID。
oozie job -oozie http://oozie-server:11000/oozie -coordinfo <coord_job_id>
将 <coord_job_id> 替换为实际的协调器 Job ID。
访问 Oozie Web UI 地址(通常为 http://oozie-server:11000/oozie),在 Web UI 上可以查看工作流和协调器的状态、日志信息、执行历史等。
5. 管理工作流或协调器:
可以使用 Oozie 命令行客户端管理工作流和协调器,例如杀死工作流、暂停协调器、恢复协调器等。
oozie job -oozie http://oozie-server:11000/oozie -kill <job_id>
oozie job -oozie http://oozie-server:11000/oozie -suspend <coord_job_id>
oozie job -oozie http://oozie-server:11000/oozie -resume <coord_job_id>