5.5 工作流调度 Oozie 本节摘要:Oozie 把前几章的组件动作(Sqoop 导入、Hive 计算、HBase 写入、MR 作业)编排为有向无环工作流,并用协调器实现定时与数据依赖触发。本节用一个日报管道的完整 workflow.xml 讲解节点模型、EL 表达式、失败重试,以及协调器的数据可用性触发语义。 为什么需要编排器 到上一节为止,数据旅程的每一段都有人开车,但没人排班。一个日报管道的真实依赖链是:凌晨两点 Sqoop 拉数 → 质量校验 MR → Hive 建分区跑聚合 → 结果导出回关系库 → 任一环节失败要告警并支持补跑。用 cron 逐个定时的问题:上游延迟十分钟,下游照点开跑,读半份数据算出错报表——定时触发不等于数据就绪。
本节摘要:Oozie 把前几章的组件动作(Sqoop 导入、Hive 计算、HBase 写入、MR 作业)编排为有向无环工作流,并用协调器实现定时与数据依赖触发。本节用一个日报管道的完整 workflow.xml 讲解节点模型、EL 表达式、失败重试,以及协调器的数据可用性触发语义。
到上一节为止,数据旅程的每一段都有人开车,但没人排班。一个日报管道的真实依赖链是:凌晨两点 Sqoop 拉数 → 质量校验 MR → Hive 建分区跑聚合 → 结果导出回关系库 → 任一环节失败要告警并支持补跑。用 cron 逐个定时的问题:上游延迟十分钟,下游照点开跑,读半份数据算出错报表——定时触发不等于数据就绪。Oozie 的价值正在这里:它按"依赖完成"而非"墙上时钟"触发下游。
Oozie 自身是个 Web 应用(跑在 Tomcat),状态存关系库,动作的实际执行仍提交给 YARN——Oozie 是排班表,不是发动机。这与 4.1 的分层逻辑一致:编排层指挥,YARN 层干活。
工作流是 XML 描述的有向无环图,节点分两族。控制节点管流程逻辑:start/end/fork/join/decision/kill。动作节点是真正的活:hive2、sqoop、map-reduce、shell、fs(HDFS 文件操作)、email。动作结束后有 ok-to 与 error-to 两条出边——每个动作都必须显式接错误路径,否则校验不过,这是 Oozie 强迫你思考失败路径的方式。
一个日报管道(节选但结构完整):
<workflow-app xmlns="uri:oozie:workflow:1.0" name="daily_report"> <start to="import"/> <action name="import"> <sqoop xmlns="uri:oozie:sqoop-action:1.0"> <resource-manager>${rmHost}</resource-manager> <name-node>hdfs://nn:8020</name-node> <arg>import</arg> <arg>--connect;jdbc:mysql://db/orders</arg> <arg>--table;t_order</arg> <arg>--target-dir;${wf:conf('inputDir')}</arg> </sqoop> <ok to="quality"/> <error to="alert"/> </action> <action name="quality"> <hive2 xmlns="uri:oozie:hive2-action:1.0"> <jdbc-url>jdbc:hive2://hs2:10000/default</jdbc-url> <script>ql/quality_check.hql</script> <param>dt=${wf:conf('dt')}</param> </hive2> <ok to="decision"/> <error to="alert"/> </action> <decision name="decision"> <switch> <case to="aggregate">${wf:actionData('quality').dirty_ratio le '0.01'}</case> <default to="quarantine"/> <!-- 脏数据超1% 进隔离分支而非直接告警死 --> </switch> </decision> <action name="aggregate"> <hive2 xmlns="uri:oozie:hive2-action:1.0"> <jdbc-url>jdbc:hive2://hs2:10000/default</jdbc-url> <script>ql/aggregate.hql</script> <param>dt=${wf:conf('dt')}</param> </hive2> <ok to="end"/> <error to="retry"/> </action> <kill name="alert"> <message>管道失败于 ${wf:lastErrorNode()}:${wf:errorMessage()}</message> </kill> <end name="end"/> </workflow-app>
两个语言细节值得点破。EL 表达式(wf:conf、wf:actionData、wf:lastErrorNode)让工作流读到参数、上一步输出与错误上下文,decision 节点靠它实现条件分支——例如质量校验动作输出 dirty_ratio,超过阈值走隔离分支。参数化:所有日期、路径用变量,同一份 XML 服务 365 天,由提交时传入的 dt 决定一切。
失败处理有两层。节点级重试:动作里加 <retry> 策略(指数退避、最多 N 次),对付偶发的连接抖动。工作流级重跑:oozie job -rerun 可指定从某节点开始(-action aggregate),已成功的导入不重做——重跑语义与第 3 章"Map 结果保留复用"如出一辙,核心理念都是幂等分段。
workflow 是一次性图,coordinator 把它变成周期服务。两种触发语义:
<coordinator-app name="daily_report_c" frequency="${coord:days(1)}" start="2026-08-01T02:00Z" end="2027-08-01T00:00Z" timezone="Asia/Shanghai" xmlns="uri:oozie:coordinator:1.0"> <datasets> <dataset name="order_input" frequency="${coord:days(1)}" initial-instance="2026-08-01T02:00Z"> <uri-template>hdfs://nn:8020/data/orders/dt=${YEAR}-${MONTH}-${DAY}</uri-template> <done-flag>_SUCCESS</done-flag> <!-- 数据就绪标记文件 --> </dataset> </datasets> <input-events> <data-in name="in" dataset="order_input"> <instance>${coord:current(0)}</instance> </data-in> </input-events> <action> <workflow><app-path>hdfs://nn:8020/apps/daily_report</app-path> <configuration> <property><name>dt</name><value>${coord:formatTime(coord:dateOffset(...),'yyyy-MM-dd')}</value></property> </configuration> </workflow> </action> </coordinator-app>
关键在 <done-flag>_SUCCESS</done-flag> 与 input-events 的组合:coordinator 到点后检查 dt 当天目录是否存在 _SUCCESS 标记文件,存在才触发;不存在就等(直到 timeout)。这就是"数据就绪触发"的实现——上游无论几点写完(手工补数、重跑都行),标记住,下游自动跟进。_SUCCESS 标记约定是整个 Hadoop 生态的通用握手信号:5.4 节 Sqoop 完成后应创建它,Hive 分区加载前应检查它,一以贯之。
bundle 把多个 coordinator 打包成产品级单元(如"销售日报产品线"含订单、流量、库存三条链),统一启停。日常运维三板斧:
oozie job -config job.properties -run # 提交workflow oozie job -run -config coord.properties # 提交coordinator oozie jobs -filter type=COORDINATOR # 看所有周期任务 oozie job -info 0000045-260818-coord # 单任务状态 oozie job -log 0000045-260818-wf # 查日志定位失败节点 oozie job -rerun 0000045-260818-wf -action 3 # 从第3个动作重跑
排障经验一条:coordinator 状态 WAITING 或 READY 卡住,九成是输入数据集的 _SUCCESS 没出现——先查上游目录,再看时区(coordinator 用 UTC 还是本地是经典坑),最后才怀疑 Oozie 本身。
至此数据的生态消费层走完。下一章换到运维视角:怎样让这条数据旅程长期稳定运行。