Article

任务调度 Azkaban

更新于:2026-07-13

第一章:Azkaban 概述与核心概念

1.1 什么是 Azkaban?

名称说明注意事项
Azkaban 定义Azkaban 是由 LinkedIn 开发并开源的工作流调度系统,用于管理和调度 Hadoop 作业流。它通过有向无环图(DAG)组织任务依赖关系,支持多种作业类型(如 Shell、Hive、Pig 等),具备 Web UI、权限控制、错误通知和调度能力。- 不是通用编程框架,而是专注于批处理任务调度。
- 强调简单性与稳定性,适合企业级 ETL 流程管理。
开源状态Apache 2.0 许可证开源项目,托管于 GitHub。社区活跃度中等,更新频率较低,但稳定性高。
核心理念”配置即代码”:通过 .job 文件定义任务,通过依赖关系构建工作流。所有逻辑通过文本文件描述,便于版本控制(如 Git)。

1.2 Azkaban 的设计目标与应用场景

名称说明注意事项
设计目标 - 可靠性确保任务在失败后可重试、可恢复,支持断点续跑机制。需配合外部存储(如 HDFS)保障数据一致性。
设计目标 - 易用性提供直观的 Web 界面,用户无需编码即可创建和监控工作流。适合非开发人员(如数据分析师)参与调度管理。
设计目标 - 扩展性支持多执行器集群部署,可横向扩展处理能力。集群模式需依赖 MySQL 统一元数据管理。
设计目标 - 可追踪性每次执行生成唯一 ID,支持日志查看、状态追踪和性能分析。日志默认保留有限时间,建议外挂日志系统。
应用场景 - ETL 流程调度常用于数据仓库每日定时抽取、清洗、加载任务。适合周期性批处理,不适用于实时流处理。
应用场景 - 数据质量检查在数据加载后自动触发校验脚本。可结合邮件或 Webhook 发送告警。
应用场景 - 模型训练调度定时启动机器学习模型训练任务(如 Spark MLlib)。需确保资源充足,避免阻塞关键流程。

1.3 Azkaban 与其他调度工具对比(如 Airflow、Oozie)

对比维度AzkabanApache AirflowApache Oozie
编程模型配置驱动(.job 文件)代码驱动(Python DAG)XML 配置文件
学习曲线低,适合初学者较高,需掌握 Python 和 DAG 概念高,XML 复杂且调试困难
UI 友好度简洁直观,重点在执行监控功能丰富,支持 Gantt 图、Task Instance 查看一般,界面较陈旧
调度灵活性支持 Cron 调度,但动态调度能力弱支持复杂调度逻辑(如 timedelta、sensor)支持 Coordinator 和 Bundle
扩展性支持多 Executor 集群可通过 Celery/K8s 扩展依赖 Hadoop 生态,扩展性一般
社区活跃度中等,更新较慢高,Apache 顶级项目,持续迭代低,Hadoop 生态衰退影响使用
适用场景中小型企业 ETL 调度大型企业复杂工作流、AI 工程化Hadoop 原生生态内任务调度
注意事项- 适合轻量级、稳定调度需求
- 不支持动态参数生成
- 运维复杂度高
- 需维护元数据库和 Scheduler
- 仅适用于 Hadoop 技术栈
- 配置繁琐

1.4 核心组件解析:Web Server、Executor Server、MySQL、H2

组件说明注意事项
Web Server接收用户请求,提供 Web UI,处理项目上传、执行触发、权限验证等。- 必须能访问 MySQL。
- 单点部署存在故障风险,建议高可用。
Executor Server实际执行 Job 的服务,从数据库拉取任务并运行,上报状态。- 可部署多个形成集群。
- 每个 Executor 需配置唯一 host:port。
MySQL存储元数据:项目信息、执行记录、调度配置、用户权限等。- 生产环境必须使用 MySQL。
- 需定期备份 schema azkaban。
H2 Database内嵌数据库,仅用于 Solo 模式测试。- 不支持多节点共享。
- 数据易丢失,禁止用于生产环境。
组件交互流程用户 → Web Server → 写入 MySQL → Executor 轮询 MySQL → 执行 Job → 回写状态- 所有通信基于数据库中转。
- Executor 与 Web Server 可跨网络部署。

1.5 工作流(Workflow)与项目(Project)基本概念

概念说明注意事项
Project(项目)一组相关 Job 的集合,是权限管理和上传的基本单位。- 必须先创建项目才能上传作业。
- 支持版本控制(每次上传为新版本)。
Job最小执行单元,对应一个 .job 文件,定义一条命令或任务。- 文件名即 Job 名称。
- 必须指定 type 属性。
Workflow(工作流)由多个 Job 按照依赖关系组成的 DAG 图,表示执行顺序。- 依赖通过 dependencies 属性定义。
- 支持并行执行无依赖节点。
.job 文件示例结构type=command
command=echo "Hello"
description=Print hello message
- 每行一个 key=value
- 不支持嵌套结构
依赖定义语法jobB:
  type=command
  command=echo "B"
  dependencies=jobA
- 多依赖用逗号分隔:dependencies=jobA,jobC
- 循环依赖会报错
Flow ID工作流的唯一标识,通常与起始 Job 同名或自定义。- 在 Web UI 中显示为工作流名称
- 可包含多个 Job 节点

第二章:Azkaban 部署与环境搭建

2.1 单机模式(Solo Server)安装与启动

方法/步骤语法/命令用途示例注意事项
下载 Azkaban 发行包wget https://github.com/azkaban/azkaban/releases/download/v3.84.4/azkaban-solo-server-0.1.0-SNAPSHOT.tar.gz获取 Solo 模式安装包(略)建议选择稳定版本(如 v3.84.4)
解压安装包tar -zxvf azkaban-solo-server-*.tar.gz展开目录结构(略)目录包含 binlibplugins 等子目录
启动服务bin/start-solo.sh启动内置 Web + H2 + Executor(无参数)自动打开 http://localhost:8081
停止服务bin/shutdown-solo.sh安全关闭服务(无参数)避免直接 kill 进程
访问 Web UI浏览器访问 http://<host>:8081登录管理界面默认账号:azkaban / azkabanH2 数据存储在本地,重启可能丢失

2.2 双服务模式(Two-Server Mode)部署

方法/步骤语法/命令用途示例注意事项
准备 Web Server解压 azkaban-web-server-*部署 Web 服务(略)需配置指向 MySQL 和 Executor 地址
准备 Executor Server解压 azkaban-executor-server-*部署执行服务(略)每个 Executor 需独立端口(默认 12321)
配置 database.propertiesdriver=com.mysql.jdbc.Driver
user=azkaban
password=azkaban
url=jdbc:mysql://localhost:3306/azkaban
Web 和 Executor 共用数据库配置放置于 conf/ 目录下必须提前创建数据库和用户
初始化数据库mysql -u root -p < scripts/create-all-sql-3.84.4.sql创建表结构替换脚本版本号脚本位于 azkaban-db 模块
启动 Web Serverbin/start-web.sh启动 Web 服务(无参数)依赖 Executor 已注册
启动 Executor Serverbin/start-executor.sh启动执行器(无参数)启动后自动向 Web 注册
停止服务bin/shutdown-web.sh / bin/shutdown-executor.sh安全关闭(无参数)建议先停 Web,再停 Executor

2.3 多执行器集群(Multi-Executor Cluster)配置

方法/步骤语法/命令用途示例注意事项
部署多个 Executor在不同机器或端口部署多个 executor 实例提升并发执行能力executor.port=12322(第二个实例)每个 executor 必须有唯一 host:port
自动注册机制Executor 启动时向 Web Server 注册实现动态发现无需手动配置Web Server 需能访问所有 Executor
负载均衡策略Web Server 轮询选择可用 Executor分摊执行压力内部实现,无需配置不支持权重分配
查看 Executor 状态Web UI → Admin → Executors监控健康状态显示 active、idle 数量失联 Executor 会标记为 DISABLED
手动禁用 Executor在 Admin 页面点击 Disable排除故障节点用于维护升级禁用后不再分配新任务
高可用保障多 Executor + Web HA(Nginx)避免单点故障前端使用负载均衡器Web Server 本身无状态,易扩展

2.4 数据库配置(MySQL)与初始化

方法/步骤语法/命令用途示例注意事项
创建数据库CREATE DATABASE azkaban DEFAULT CHARACTER SET utf8;初始化存储空间(SQL 语句)建议使用 utf8 编码
创建用户CREATE USER 'azkaban'@'%' IDENTIFIED BY 'azkaban';
GRANT ALL PRIVILEGES ON azkaban.* TO 'azkaban'@'%';
授权访问可限制 IP 范围生产环境避免使用 %
配置 database.propertiesdriver=com.mysql.jdbc.Driver
user=azkaban
password=azkaban
url=jdbc:mysql://dbhost:3306/azkaban?useSSL=false
连接 MySQL放于 Web 和 Executor 的 conf/ 目录useSSL=false 可避免证书问题
导入表结构mysql -u azkaban -p azkaban < create-all-sql-*.sql创建元数据表脚本来自源码或发行包确保版本匹配
测试连接bin/start-web.sh 观察日志验证数据库可达性grep "connected" logs/失败常见于驱动缺失或网络不通
JDBC 驱动放置mysql-connector-java-x.x.x.jar 放入 lib/ 目录确保 JDBC 可用(文件操作)Web 和 Executor 都需要

2.5 常见部署问题排查

问题现象可能原因解决方法注意事项
Web 页面无法访问端口未开放、服务未启动检查 8081 端口 `netstat -tlnpgrep 8081`
Executor 未注册network.timeout 设置过小增大 azkaban.properties 中的 network.timeout默认 5000ms,可设为 10000
数据库连接失败驱动缺失、URL 错误、权限不足检查 database.properties 和 JDBC 驱动使用 telnet 测试数据库连通性
上传项目失败磁盘空间不足、权限不足检查 executor 本地工作目录(如 executions/默认路径为安装目录下
工作流卡在 QUEUED 状态所有 Executor 忙或失联查看 Admin → Executors 状态可重启 Executor 或增加实例
中文乱码JVM 未设置 UTF-8启动脚本添加 -Dfile.encoding=UTF-8start-web.sh 中修改 JAVA_OPTS
时间不同步节点间系统时间差异大使用 NTP 同步所有服务器时间时间偏差可能导致调度异常

第三章:创建与管理 Azkaban 项目

3.1 使用 Web UI 创建项目

方法语法/操作步骤用途示例注意事项
登录 Web UI浏览器访问 http://<azkaban-web>:8081,输入用户名和密码(如 azkaban/azkaban进入主界面默认账户由 azkaban-users.xml 定义确保网络可达且服务已启动
创建新项目点击 “Create Project”,填写 Project Name、Description、Owner初始化一个空项目容器Name: etl_daily
Description: 每日ETL任务
项目名不可重复,不支持特殊字符
填写项目属性设置项目描述、负责人、权限组便于团队协作与审计Owner: data_team描述建议清晰说明用途
提交创建点击 “Create” 按钮在数据库中持久化项目元数据(无返回值)成功后跳转至项目主页

3.2 使用命令行工具(azkaban-cli)上传项目

方法语法/命令用途示例注意事项
打包项目目录zip -r project.zip *.job *.properties将 job 文件打包为 zip包含所有依赖脚本和配置文件必须为 zip 格式,不支持 tar/gz
获取 Session IDcurl -k -X POST --data 'action=login&username=azkaban&password=azkaban' https://localhost:8443认证并获取 session返回 JSON 中包含 session.id-k 忽略 SSL 证书验证
上传并覆盖项目curl -k -X POST -H "Content-Type: multipart/form-data" --form project=project_name --form file=@project.zip --form ajax=upload --form session.id=<session_id> https://localhost:8443/manager上传项目到服务器替换 project_namesession.id若项目不存在则创建,存在则覆盖
查看上传结果命令返回 JSON 响应确认是否成功{ "status": "success", "projectId": 12 }失败时检查 session 是否过期
注意事项- 上传前需确保项目已通过 Web 创建(除非使用 API 创建)
- ZIP 包大小受限于 max.request.size 配置
- 推荐用于 CI/CD 自动化部署

3.3 项目权限管理(Permissions & Roles)

角色权限说明可执行操作注意事项
Admin(管理员)全部权限创建/删除项目、上传作业、执行工作流、修改权限、查看日志可修改他人作业,慎用
Read(只读)查看权限查看项目结构、执行历史、日志输出无法触发或修改任何内容
Write(写入)编辑权限上传新版本、修改 job 文件(需配合 Execute)不包含执行权限
Execute(执行)运行权限手动启动工作流、重试节点需配合 Read 查看结果
设置权限方式Web UI → Project → Permissions → 添加用户并选择角色分配用户权限支持通配符组(如 LDAP group)
默认权限策略新建项目仅 Owner 拥有 Admin 权限安全隔离其他用户需显式授权

3.4 项目版本控制与更新策略

概念说明示例注意事项
版本机制每次上传 zip 包视为一个新版本Web UI 显示版本号(v1, v2…)可回滚到任意历史版本
查看历史版本Project → Executions → Versions Tab审计变更记录支持下载旧版 zip 包
回滚操作选择某一历史版本 → Click “Rollback”恢复至稳定状态适用于错误发布后修复
版本命名建议结合 Git 提交 ID 或时间戳命名 zip 文件etl_v20250315_build_abc123.zip便于追踪来源
并行开发策略开发分支独立打包测试 → 合并后上传主项目多人协作场景避免直接在生产项目修改

3.5 删除与归档项目

操作方法/命令用途示例注意事项
删除项目Web UI → Project → Settings → Delete Project永久移除项目及其所有执行记录输入项目名确认删除无法恢复,请谨慎操作
归档项目(软删除)重命名项目为 [ARCHIVED]_old_name,移除所有用户权限标记废弃但保留数据[ARCHIVED]_temp_job_2024便于后续审计或恢复
批量清理脚本编写 Python 脚本调用 REST API 批量删除自动化运维遍历 /manager 接口查询项目列表建议先导出元数据备份
数据库级清理DELETE FROM projects WHERE name LIKE 'temp%';
DELETE FROM execution_flows WHERE project_id = ?;
彻底清除数据需同步删除 executions 表操作前必须备份数据库
注意事项- 删除项目不会自动清理 HDFS 或本地临时文件
- 建议先停止所有调度任务再删除

第四章:编写 Azkaban 工作流(.flow 文件)

4.1 job 文件基础语法结构

属性语法格式用途示例注意事项
typetype=command | java | hive | flow定义任务类型type=command必须指定,决定执行器行为
commandcommand=echo "hello"执行 Shell 命令command=hive -f /path/to/script.hql多行命令可用 \ 换行
descriptiondescription=This is a test job添加注释说明提高可读性Web UI 中显示
retriesretries=3失败自动重试次数每次间隔由 retry.backoff 值决定不适用于关键一致性任务
retry.backoffretry.backoff=10000重试间隔(毫秒)默认 10 秒避免频繁重试压垮系统
memory / cpumemory.mb=2048
cpu.threads=2
资源限制(需插件支持)控制作业资源占用非强制,依赖外部监控
注释# 这是一行注释忽略该行内容支持行尾注释不支持多行 /* */

4.2 定义节点依赖关系(dependencies)

方法语法用途示例注意事项
单依赖dependencies=jobA当 jobA 成功后执行当前任务jobB:
  type=command
  command=echo "B"
  dependencies=jobA
依赖任务必须在同一项目中
多依赖dependencies=jobA,jobC所有前置任务成功后才运行jobD:
  type=command
  command=merge_data.sh
  dependencies=extract1,extract2
逗号分隔,无空格
并行执行多个任务无依赖关系同时启动多个分支jobA → jobC
jobB → jobC
提升整体执行效率
循环依赖检测系统自动检测防止死锁jobA → jobB → jobA提交时报错 “Circular dependency”
依赖命名规范使用有意义的名称提高可维护性extract_user_data, transform_sales避免使用 job1, job2

4.3 使用内联工作流(Inline Flow in job file)

方法语法用途示例注意事项
定义内联 Flowtype=flow
flow.name=myFlow
nodes=[{...}]
在 job 文件中直接定义复杂 DAG见下方示例适合小型嵌套流程
nodes 数组结构nodes=[{"name":"A","type":"command","command":"echo A"}, {"name":"B","dependencies":["A"],"command":"echo B"}]描述 DAG 节点JSON 格式嵌入注意引号转义
嵌套 Flow 调用type=embeddedFlow
flow.project=project_x
flow.flowId=flow_y
调用其他项目的 Flow实现模块复用调用方需有 Read 权限
示例:简单内联 Flowmyflow.job:
type=flow
flow.name=inline_test
nodes=[{"name":"start","type":"command","command":"echo start"},{"name":"end","dependencies":["start"],"type":"command","command":"echo done"}]
替代多个 .job 文件减少文件数量复杂逻辑仍建议拆分
注意事项- JSON 必须合法,避免语法错误
- 不支持跨项目内联定义

4.4 复杂 DAG 图构建技巧

技巧说明示例注意事项
分层设计按阶段划分:Extract → Transform → Load每层多个并行任务汇聚到下一阶段降低耦合度,便于调试
子流程复用将通用逻辑封装为独立 Flowdaily_cleanup.flow通过 embeddedFlow 调用
动态参数传递使用 ${input_path} 在运行时替换配合执行时传参实现灵活调度参数需提前声明或传入
错峰执行设置高耗时任务错开时间避免资源争抢使用调度延迟或优先级
虚拟节点(Dummy Job)添加空任务作为汇合点wait_all:
  type=noop
  dependencies=task1,task2
noop 类型不执行实际操作

4.5 错误处理与失败策略(fail-action, on-failure)

属性语法用途示例注意事项
fail-actionfail-action=finishCurrent定义失败后的整体行为finishCurrent | cancelImmediately | retry作用于整个 Flow
on-failureon.failure=finishCurrent同 fail-action,旧写法推荐统一使用 fail-action两者功能相同,优先用新语法
finishCurrentfail-action=finishCurrent当前任务失败后,允许已完成任务继续,不启动新任务适用于非关键路径任务最终 Flow 状态为 FAILED
cancelImmediatelyfail-action=cancelImmediately一旦失败立即终止所有运行和待运行任务关键任务链使用防止脏数据传播
retry(节点级)retries=3
retry.backoff=5000
单个任务自动重试适用于瞬时故障(如网络抖动)不适用于数据一致性破坏场景
自定义失败处理 Jobjob_failure_handler:
  type=command
  command=notify_error.sh ${azkaban.flow.execid}
  on.failure=END
失败时执行通知脚本发送邮件或 Webhook需确保通知服务可靠

第五章:Job 类型与处理器(Job Types)

5.1 Command Job Type(command)

方法/属性语法用途示例注意事项
typetype=command定义为 Shell 命令任务type=command必须指定
commandcommand=<shell_command>执行任意 Shell 命令command=echo "Hello"
command=sh /data/scripts/cleanup.sh
支持多行(用 \ 连接)
working.dirworking.dir=/path/to/dir设置工作目录working.dir=/home/azkaban/jobs默认为执行器临时目录
env.env.PATH=/usr/local/bin:$PATH设置环境变量env.HADOOP_HOME=/opt/hadoop仅对该 Job 有效
ulimitulimit=-v 8388608限制内存使用(KB)ulimit=-u 64防止资源耗尽
注意事项- 不支持交互式命令
- 建议将脚本外部化而非写在 job 文件中

5.2 Java Job Type(java)

方法/属性语法用途示例注意事项
typetype=java调用本地 JVM 执行 Java 类type=java类必须可被 classpath 加载
classclass=com.example.MainClass指定主类名class=com.etl.BatchProcessor必须包含 main 方法
jarsjars=/path/to/lib/a.jar,/path/to/lib/b.jar添加依赖 JAR 包路径jars=/home/azkaban/lib/utils.jar多个用逗号分隔
argsargs=arg1 arg2传入 main 方法参数args=input.parquet output.orc空格分隔多个参数
jvm.argsjvm.args=-Xmx2g -Dlog.level=INFO设置 JVM 启动参数jvm.args=-server -XX:+UseG1GC影响性能和稳定性
lib.dirlib.dir=/path/to/libs自动加载该目录下所有 jarlib.dir=/home/azkaban/java_libs简化依赖管理
注意事项- 主类需打包进 jar 并确保无冲突依赖
- 推荐使用 hadoopJava 处理 Hadoop 作业

5.3 Pig Job Type(pig)

方法/属性语法用途示例注意事项
typetype=pig使用本地 Pig 引擎运行脚本type=pig已废弃,建议用 hadoopPig
pig.scriptpig.script=/path/to/script.pig指定 Pig 脚本路径pig.script=/home/user/analyze.pig脚本必须存在且可读
parametersparameters="input=data.csv output=result"传入参数parameters="year=2025 debug=true"双引号包裹多个参数
注意事项- 该类型不连接 Hadoop 集群
- 仅用于测试或单机模式

5.4 Hadoop Pig Job Type(hadoopPig)

方法/属性语法用途示例注意事项
typetype=hadoopPig在 Hadoop 集群上运行 Pig 脚本type=hadoopPig需配置 Hadoop 环境
pig.scriptpig.script=/path/to/script.pig指定 Pig 脚本pig.script=hdfs://nn:9000/pig/etl.pig支持 HDFS 路径
hadoop.security.manager.classhadoop.security.manager.class=org.apache.azkaban.jobtype.HadoopSecurityManager_H_2_0指定安全管理器根据 Hadoop 版本选择Azkaban 插件需支持
job.flow.idjob.flow.id=${azkaban.flow.execid}传递 Flow ID(可选)用于日志追踪通常自动注入
parametersparameters="in=hdfs://... out=hdfs://..."传入 Pig 参数parameters="date=${today}"支持变量替换
注意事项- Executor 必须安装 Pig 并配置好 Hadoop 客户端
- 脚本中可用 $parameter 接收参数

5.5 Hadoop WordCount Job Type(hadoopWordCount)

方法/属性语法用途示例注意事项
typetype=hadoopWordCount简化版 MapReduce 任务模板type=hadoopWordCount专为 WordCount 示例设计
input.pathinput.path=/data/input.txt输入文件路径input.path=hdfs://nn:9000/input支持通配符如 *.txt
output.pathoutput.path=/data/output输出目录(必须不存在)output.path=hdfs://nn:9000/out/2025任务失败前会自动删除
hadoop.job.classhadoop.job.class=com.hadoop.mapreduce.WordCount自定义 MapReduce 类替换默认实现必须继承相应接口
注意事项- 实际生产中不推荐使用此类型
- 更通用的做法是使用 hadoopJava 或直接提交 Jar

5.6 Hive Job Type(hive)

方法/属性语法用途示例注意事项
typetype=hive使用本地 Hive CLI 执行 HQLtype=hive不连接 HiveServer2,已过时
hive.scripthive.script=/path/to/query.hql指定 HQL 脚本路径hive.script=/home/azkaban/sql/dim_user.hql文件需 UTF-8 编码
hive.homehive.home=/opt/hive设置 Hive 安装路径hive.home=/usr/lib/hive影响 bin/hive 调用
hive.parametershive.parameters="dt=2025-01-01 env=prod"传入 Hive 变量SELECT * FROM logs WHERE dt='${dt}'在 HQL 中用 ${} 引用
注意事项- 建议改用 hadoopHive 类型以支持远程执行
- 本地模式需安装 Hive 客户端

5.7 EmbeddedFlow Job Type(嵌套流程)

方法/属性语法用途示例注意事项
typetype=embeddedFlow调用当前项目或其他项目的子流程type=embeddedFlow实现模块化设计
flow.projectflow.project=my_subproject指定目标项目名flow.project=common_libs省略则表示当前项目
flow.flowIdflow.flowId=cleanup_flow指定要调用的工作流 IDflow.flowId=daily_export必须存在且可访问
propagate.failurepropagate.failure=true是否传播子流程失败true | false设为 true 则父流程也失败
flow.parametersflow.parameters="param1=value1,param2=value2"向子流程传参flow.parameters="date=${today}"支持变量替换

示例:

# subflow.job
type=embeddedFlow
flow.project=utils
flow.flowId=vacuum_db
flow.parameters="days=7"

子项目需授予 Read 权限。

5.8 自定义 Job Type 扩展机制

方法/属性说明示例注意事项
创建插件目录在 executor 的 plugins/ 目录下创建 jobtype/my-custom-job/结构:
- plugin.properties
- lib/*.jar
- classes/
必须符合 Azkaban 插件规范
plugin.propertiesname=myCustomType
class=com.example.MyJobRunner
注册新 Job 类型name 是 job 文件中使用的 type 值
实现 JobRunnerpublic class MyJobRunner extends AbstractJob { ... }编写执行逻辑覆盖 run() 方法
打包部署将编译后的 class 和依赖打成 jar 放入 lib/mvn package 或手动打包重启 Executor 生效
配置白名单azkaban.properties 中添加:jobtype.plugin.dirs=my-custom-job允许加载自定义类型多个用逗号分隔
注意事项- 需重新编译并部署 Executor
- 错误可能导致整个 Executor 异常
- 建议先在测试环境验证

第六章:参数传递与变量替换

6.1 内置变量(如 ${azkaban.flow.execid}

变量名说明示例值注意事项
${azkaban.job.id}当前 Job 的唯一标识extract_users用于日志标记
${azkaban.flow.id}当前工作流的 IDdaily_etl_flow可用于路径组织
${azkaban.flow.execid}当前执行实例的编号12345每次运行不同,适合临时目录
${azkaban.executor.host}执行该 Job 的 Executor 主机名node01.cluster用于分布式调试
${azkaban.project.version}项目当前版本号v17便于审计变更
${azkaban.project.dir}项目解压后的本地路径/tmp/azkaban/executions/12345可访问上传的脚本文件
注意事项- 所有变量在运行时自动替换
- 不区分大小写?否,必须原样引用

6.2 用户自定义参数传入方式

方法语法用途示例注意事项
job 文件内定义myparam=value静态参数input_path=/data/raw修改需重新上传
Web UI 手动输入在执行页面添加 key=value临时覆盖date=2025-09-30仅本次生效
命令行上传时传参--params "k1=v1&k2=v2"CI/CD 场景自动化curl ...&date=${DATE}URL 编码特殊字符
调度配置中设置在 Schedule 页面填写 Parameters定时任务固定参数interval=1d每次调度都使用
注意事项- 参数优先级:运行时 > 调度 > job 文件默认值
- 不支持嵌套表达式

6.3 运行时参数覆盖(Runtime Properties)

方法语法用途示例注意事项
执行时传参在 Execute Now 对话框中添加键值对动态控制行为dry_run=true覆盖所有层级默认值
REST API 传参"overrideProps": {"key":"value"}程序化调用用于自动化触发JSON 格式
参数文件--props runtime.properties批量传入参数文件内容:lineage=off适合复杂配置
作用域仅影响本次执行不修改项目元数据
注意事项- 参数名不能以 azkaban. 开头(保留字)
- 无法覆盖 type, dependencies 等结构属性

6.4 跨 Job 参数共享(通过 flow.parameters)

方法语法用途示例注意事项
定义 flow.parametersflow.parameters=["param1", "param2"]声明流程级共享参数.job.flow 文件中定义必须提前声明才能传递
在 Job 中引用command=process.sh ${myparam}使用变量myparam 必须在 flow.parameters否则替换为空字符串
嵌套流程传参flow.parameters="shared_param=${top_level_value}"向子流程传递实现上下文透传需双方约定参数名

示例结构:

# parent.job
type=flow
flow.name=main
flow.parameters=[ "date" ]
nodes=[{...}]

# childA.job
type=command
command=echo ${date}

所有节点均可访问 date,实现全局日期控制。

| 注意事项 | - flow.parameters 是显式声明机制
- 非声明式变量不会自动传播 | | | |

6.5 安全参数处理(secure properties)

方法语法用途示例注意事项
secure.props.filesecure.props.file=/path/to/secure.properties指定加密属性文件路径放置于 executor conf/ 目录Web Server 不需要
secure.property.namessecure.property.names=db.password,api.key声明哪些参数为敏感信息多个用逗号分隔必须匹配实际参数名
传参方式在执行时传入 db.password=xxx但日志和 UI 中隐藏显示为 ******防止密码泄露
属性文件格式db.password=encrypted(AES):base64data
api.key=plain:text_key
支持加密或明文推荐使用加密存储工具类可生成密文
注意事项- 仅 Executor 能解密读取
- Web UI 和日志均不显示原始值
- 必须保证所有 Executor 配置一致

第七章:触发器与调度机制

7.1 手动执行工作流

方法语法/操作步骤用途示例注意事项
Web UI 执行进入项目 → 点击工作流 → “Execute Flow” → 填写参数 → Submit临时触发一次执行如修复数据后手动运行 ETL不影响定时调度计划
REST API 触发POST /executor?
ajax=executeFlow
&project=<project_name>
&flow=<flow_id>
&session.id=
程序化调用执行结合脚本或监控系统自动触发需先获取有效 session.id
传入运行时参数在 Execute 页面添加 key=value 对覆盖默认配置date=2025-09-30
dry_run=true
支持变量替换
查看执行结果跳转至 Execution 页面,观察状态变化实时监控执行过程成功 → SUCCEEDED
失败 → FAILED
可点击查看节点日志
注意事项- 手动执行生成独立 execid
- 若已存在定时任务,不影响其后续运行
- 支持并行执行多个实例

7.2 定时调度(Scheduling)配置

方法操作步骤用途示例注意事项
启用调度工作流页面 → Schedule → 设置时间与频率创建周期性任务每日凌晨 2 点执行需项目具有 Execute 权限
设置首次执行时间指定 “First Run Time”控制启动时机2025-10-01 02:00可延迟启动
设置重复频率选择间隔:Minutes / Hours / Days / Weeks定义调度周期Daily at 02:00 AM支持自定义 Cron
添加运行参数在调度配置中填写 key=value固定每次调度的输入env=prod
batch_size=10000
参数对所有调度实例生效
保存调度点击 “Schedule” 按钮提交调度任务到数据库成功后显示下一次执行时间调度信息存储在 MySQL 中
注意事项- 调度由 Web Server 统一管理
- Executor 不参与调度决策
- 修改调度需重新提交

7.3 Cron 表达式在 Azkaban 中的应用

字段位置含义允许值示例说明
10–590Azkaban 默认从分钟开始,秒固定为 0
2分钟0–590每小时整点触发
3小时0–232凌晨 2 点
4日期1–31*每天
5月份1–12*每月
6星期0–6(0=Sunday)1–5周一至周五
Cron 表达式示例用途对应配置注意事项
0 0 2 * * ?每天凌晨 2:00 执行标准 daily 调度Azkaban 使用 Quartz 兼容格式,但秒位固定为 0
0 0 0/1 * * ?每小时整点执行用于实时性要求高的任务高频调度注意资源消耗
0 30 8 ? * MON-FRI工作日早上 8:30 执行业务上班前完成数据准备? 表示不指定具体日期
0 0 0 1 * ?每月 1 日午夜执行月度报表生成注意跨月边界

| 注意事项 | - Azkaban 的 Cron 不支持年字段
- 推荐使用 Web UI 的图形化调度器生成表达式
- 错误表达式会导致调度失败且无明显提示 | | | |

7.4 触发器(Triggers)与条件执行(实验性功能)

概念说明示例注意事项
Trigger(触发器)实验性功能,用于基于事件或时间触发工作流如”当 HDFS 文件到达时启动”需自定义开发,官方支持有限
条件节点(Conditional Job)根据前序任务输出决定是否执行exit.status == 0 则继续需结合脚本返回值实现
外部信号触发通过 REST API 或消息队列通知 AzkabanKafka 消息 → 调用 executeFlow API实现近实时响应
Sensor 类型 Job(模拟)使用循环检查 + sleep 实现等待while [ ! -f /flag.txt ]; do sleep 60; done占用资源,不优雅
注意事项- Azkaban 本身不提供成熟的事件驱动机制
- 条件逻辑建议在 Shell 脚本中实现
- 可结合 Airflow Sensor 做高级调度

7.5 调度状态监控与修改

操作方法/路径用途示例注意事项
查看调度列表Web UI → Project → Schedules Tab显示所有已配置的定时任务包含下一次执行时间、状态状态包括 ACTIVE / PAUSED
暂停调度点击 “Pause” 按钮临时停止自动执行如系统维护期间可随时恢复
恢复调度点击 “Resume” 按钮重新激活已暂停的调度恢复后按原计划运行下次执行时间重新计算
修改调度时间编辑现有 Schedule → 调整 Cron 或频率更新执行计划从每天改为每小时修改后立即生效
删除调度点击 “Unschedule”彻底移除定时配置不影响手动执行能力无法撤销,请确认
数据库级查询SELECT * FROM schedule WHERE project_name='etl';审计或批量处理用于运维脚本不建议直接修改表数据
注意事项- 所有调度状态变更记录在 MySQL 中
- 多个 Executor 共享同一调度源
- Web Server 故障会导致调度中断

第八章:执行监控与日志分析

8.1 查看执行记录(Executions)

方法操作路径用途示例注意事项
进入执行历史项目主页 → Executions Tab查看所有运行实例按 execid 排序,最新在前支持分页
筛选执行状态使用状态过滤器:SUCCEEDED / FAILED / RUNNING快速定位问题查找最近失败的任务可结合时间范围筛选
查看执行详情点击某个 execid进入 DAG 图页面显示节点依赖与执行顺序支持缩放与展开
查看执行参数在 Execution Info 区域查看 “Flow Parameters”审计输入配置确认是否传入正确参数敏感参数显示为 ******
导出执行报告手动复制或截图用于汇报或归档包含开始/结束时间、耗时无原生导出功能,需定制
注意事项- 执行记录长期保留依赖数据库容量
- 可配置日志清理策略(如 TTL)
- 推荐结合外部监控系统告警

8.2 实时日志查看与下载

方法操作方式用途示例注意事项
查看节点日志点击 DAG 中某个节点 → “View Logs”实时观察输出日志自动刷新(轮询)延迟约 1–3 秒
日志分页支持 Previous / Next 分页浏览大量日志每页默认 100 行可调整大小
下载完整日志点击 “Download” 按钮保存到本地分析文件名为 <job_id>.log包含标准输出和错误
查看执行器本地日志登录 Executor 服务器 → 查看 logs/executions/<execid>/深度排查问题如 Web UI 日志缺失路径由 execution.log.dir 配置
日志保留策略默认保留最近 N 次执行日志防止磁盘溢出可配置 executor.log.retention.days建议外挂日志系统(如 ELK)
注意事项- 日志写入异步,可能短暂延迟
- 大日志可能导致浏览器卡顿
- 下载日志不包含变量替换过程

8.3 节点状态追踪(Running / Failed / Succeeded)

状态说明如何查看注意事项
READY等待依赖完成DAG 图中灰色节点仅当所有前置任务成功后进入 RUNNING
RUNNING正在执行节点高亮为蓝色可查看实时日志
SUCCEEDED成功完成节点变为绿色不再重试
FAILED执行失败节点变为红色显示失败原因摘要
DISABLED被禁用(条件未满足)灰色带斜线常见于 on-failure=END 的后续任务
SKIPPED跳过执行浅灰色如父流程失败且 fail-action=cancelImmediately
查看状态详情点击节点 → “View Info”显示开始/结束时间、主机、返回码return code=0 表示成功
注意事项- 状态由 Executor 上报至数据库
- 网络中断可能导致状态滞留(如长时间 RUNNING)
- 可通过 REST API 查询状态

8.4 重试失败节点(Retry Node)

方法操作步骤用途示例注意事项
重试单个节点点击失败节点 → “Retry” 按钮重新执行该任务适用于临时故障(如网络抖动)原执行记录保留,新增尝试
重试整个工作流点击 “Rerun Failed” 或 “Execute Again”重新运行全部或失败分支”Rerun Failed” 仅重试失败路径注意幂等性,避免重复写入
重试行为系统生成新 attempt_id,复用原参数保持上下文一致日志路径包含 attempt_1, attempt_2最大重试次数可配置
查看重试日志查看对应 attempt 的日志文件对比前后差异分析失败是否解决路径:/executions/<execid>/<job>_<attempt>.log
注意事项- 重试不会改变原 execid
- 需确保任务具备幂等性(如先删除再写入)
- 不建议对数据写入类任务频繁重试

8.5 性能瓶颈分析建议

分析维度方法工具/路径建议
节点耗时分析查看每个 Job 的 DurationExecution 页面时间轴找出最长耗时任务,优化 SQL 或脚本
并行度不足DAG 中存在长链无分支优化依赖结构拆分独立任务并行执行
资源竞争多个 Job 同时访问同一数据库或 HDFS 路径系统监控工具(如 Grafana)错峰调度或增加资源
Executor 负载过高某个 Executor 长时间 busyWeb UI → Admin → Executors增加 Executor 实例或负载均衡
日志写入慢日志刷新延迟严重检查磁盘 I/O 或网络使用 SSD 或优化日志级别
数据倾斜某个 MapReduce 任务远慢于其他查看 Hadoop/YARN 日志优化分区或 Key 设计
数据库压力MySQL CPU 或连接数过高show processlist / 监控面板优化查询或分库分表
通用建议- 启用压缩减少 I/O
- 使用分区表避免全表扫描
- 避免在高峰时段调度大任务
- 定期归档历史执行记录

第九章:安全机制与权限控制

9.1 用户认证方式(LDAP / JDBC / File-based)

认证方式配置方法用途示例注意事项
File-based(文件认证)azkaban-users.xml 中定义用户适用于小型环境或测试<user username="admin" password="admin" roles="admin"/>明文存储密码,不推荐生产使用
LDAP 认证设置 azkaban.properties
executor.use.executor.ldap=true
ldap.hostname=ldap.company.com
ldap.port=389
ldap.base.dn="ou=users,dc=company,dc=com"
集中管理企业用户支持 Active Directory需网络可达且证书可信
JDBC 认证配置数据库连接信息:
database.type=mysql
database.host=localhost
database.user=azkaban
database.password=***
并实现自定义 UserDAO
与现有系统集成用于已有用户中心的场景需开发插件支持
启用认证类azkaban.user.manager=azkaban.user.XmlUserManager指定用户管理器实现File: XmlUserManager
LDAP: LdapUserManager
必须重启 Web Server 生效
角色映射roles 配置决定权限级别admin 可访问所有项目可自定义 roles 文件推荐最小权限原则
注意事项- 生产环境禁用 file-based 认证
- LDAP 推荐启用 SSL(ldaps)
- 所有认证信息不加密传输时存在风险

9.2 项目级权限设置(Admin / Read / Write / Execute)

权限类型可执行操作适用角色注意事项
Admin(管理员)创建/删除项目、修改权限、上传作业、执行工作流、查看日志项目负责人、运维人员拥有最高权限,慎分配
Read(只读)查看项目结构、执行历史、日志输出开发人员、审计人员无法修改或触发任何内容
Write(写入)上传新版本 job 文件、修改配置开发人员不包含执行权限,需配合 Execute
Execute(执行)手动启动工作流、重试节点调度员、数据工程师需 Read 权限配合查看结果
设置方式Web UI → Project → Permissions → 添加用户并选择角色图形化操作支持用户名或 LDAP 组
权限继承无层级继承机制每个项目独立授权需手动同步团队权限
注意事项- 权限粒度为项目级,不支持目录或 Job 级别
- 删除用户后其创建的 Job 仍保留
- 建议按团队建立统一权限模板

9.3 SSL 配置与 HTTPS 支持

配置项语法/路径用途示例注意事项
启用 HTTPSazkaban.properties 中设置:
jetty.ssl.context.key.store.path=/path/to/keystore.jks
jetty.ssl.context.key.store.password=changeit
jetty.ssl.context.trust.store.path=/path/to/truststore.jks
加密 Web 通信使用 JKS 或 PKCS12 格式必须重启 Web Server
生成密钥库keytool -genkeypair -alias azkaban -keyalg RSA -keystore keystore.jks -storepass changeit -validity 365创建自签名证书适用于测试环境生产建议使用 CA 签名证书
配置端口jetty.ssl.port=8443HTTPS 监听端口默认关闭 HTTP(jetty.port=-1避免明文传输
客户端验证jetty.ssl.context.need.client.auth=true双向认证(mTLS)提高安全性需分发客户端证书
证书格式兼容性支持 JKS、PKCS12便于迁移keystore.type=PKCS12推荐使用现代格式
注意事项- 所有 Executor 必须信任 Web Server 证书(导入 truststore)
- 若使用反向代理(如 Nginx),可在前端终止 SSL
- 避免使用弱加密算法(如 SSLv3)

9.4 敏感信息加密存储策略

方法配置/语法用途示例注意事项
Secure Properties 机制azkaban.properties 中设置:
secure.props.file=/conf/secure.properties
secure.property.names=db.password,api.key
声明敏感参数名执行时传入但日志中隐藏仅 Executor 能解密
加密属性文件db.password=encrypted(AES):base64data
api.key=plain:text_key
存储加密后的值工具类生成密文支持 AES 加密
参数传入方式执行时通过 UI 或 API 传入 db.password=xxx运行时注入显示为 ******不记录在 job 文件中
加解密工具使用 Azkaban 提供的 CryptoUtils 类生成加密字符串java -cp azkaban-common.jar CryptoUtils encrypt password需相同密钥环境
密钥管理encryption.key=your-secret-key对称加密密钥存放于配置文件或环境变量必须保护密钥安全
注意事项- 所有 Executor 必须拥有相同的 secure.props.file 和密钥
- Web Server 不解密,仅传递
- 不支持嵌套加密参数

9.5 多租户隔离实践

隔离维度实现方式说明建议
项目命名空间按团队/部门前缀划分项目名teamA_etl_daily, teamB_report便于权限管理和审计
权限控制每个项目独立设置 Read/Write/Execute 权限禁止跨项目访问避免 Admin 权限滥用
Executor 分组使用 executor.group 标签隔离资源teamA_jobs → group=A
teamB_jobs → group=B
防止资源争抢
存储路径隔离作业脚本、日志、临时文件使用独立路径working.dir=/data/azkaban/${project}避免文件冲突
数据访问控制结合 Hive/Ranger 实现数据层权限即使 Job 可运行,也无法读取未授权数据深度防御策略
调度优先级通过自定义插件实现优先级队列关键业务任务优先调度需扩展调度器逻辑
注意事项- Azkaban 原生多租户能力有限
- 建议结合外部调度平台(如 Airflow)做统一治理
- 审计日志应记录操作主体与项目上下文

第十章:高可用与集群运维

10.1 多 Executor 负载均衡

方法配置/操作用途示例注意事项
注册多个 Executor启动多个 Executor 实例,自动向 Web Server 注册分散执行压力每台机器部署一个 Executor需网络互通
轮询调度Web Server 默认采用轮询策略分配任务均匀分发负载无需额外配置简单但不考虑实际负载
按资源调度实验性支持 CPU/Memory 标签高负载任务分配给高性能节点executor.resources.cpu=8
executor.resources.mem=16g
需自定义调度器
Executor 分组使用 executor.group=etlrealtime按任务类型隔离将批处理与实时任务分开提高稳定性
查看负载状态Web UI → Admin → Executors监控各节点活跃任务数busyCount 字段显示当前负载可识别瓶颈节点
注意事项- 所有 Executor 必须能访问共享存储(如 HDFS)
- 配置文件(如 secure.props)需保持一致
- 推荐使用 DNS 或 VIP 统一访问

10.2 Executor 故障转移机制

机制说明触发条件恢复行为注意事项
心跳检测Executor 每 30 秒向 Web Server 发送心跳心跳超时(默认 60 秒)Web Server 标记为 DISCONNECTED可配置 executor.heartbeat.interval
任务重试已分配但未完成的任务Executor 断开连接由 Web Server 重新调度到其他节点保证最终执行
状态持久化执行状态存储在 MySQL 数据库即使 Executor 宕机重启后可继续上报日志防止状态丢失
不完全失败处理若任务已开始执行但未完成Executor 挂掉无法判断是否完成,可能重复执行建议任务具备幂等性
手动干预管理员可强制标记任务为失败或跳过长时间无响应使用 Web UI 修改状态用于紧急恢复
注意事项- 故障转移不能保证 exactly-once 语义
- 大任务中途失败可能导致资源浪费
- 建议设置合理的超时与重试策略

10.3 Web Server 高可用部署(Nginx + HAProxy)

组件配置要点用途示例注意事项
HAProxyfrontend azkaban_https
bind *:8443 ssl crt /certs/azkaban.pem
default_backend azkaban_servers
backend azkaban_servers
balance roundrobin
server web1 192.168.1.10:8443 check ssl verify none
server web2 192.168.1.11:8443 check ssl verify none
TCP 层负载均衡支持 SSL 终止或透传verify none 用于自签名证书
Nginxupstream azkaban {
server 192.168.1.10:8443;
server 192.168.1.11:8443;
}
server {
listen 443 ssl;
location / { proxy_pass https://azkaban; }
}
HTTP 层反向代理可做路径路由、缓存推荐用于静态资源
共享数据库所有 Web Server 连接同一 MySQL 实例状态共享使用 RDS 或主从复制必须保证数据库高可用
会话粘性(Session Stickiness)不需要Azkaban 无用户会话状态所有操作基于数据库可自由扩展 Web Server
健康检查/index/status 接口检测后端可用性HAProxy 使用 check 参数自定义健康检查脚本
注意事项- Web Server 是无状态的,可水平扩展
- Executor 必须注册到所有 Web Server 或通过 VIP 访问
- 建议使用 Keepalived 实现 VIP 故障转移

10.4 数据库备份与恢复

操作方法工具/命令注意事项
全量备份mysqldump -u azkaban -p azkaban_db > backup.sql定期导出数据建议每日凌晨执行
增量备份启用 binlog:
log-bin=mysql-bin
server-id=1
记录所有变更结合 xtrabackup
自动化脚本编写 shell 脚本 + crontab实现周期备份gzip 压缩节省空间
恢复流程mysql -u azkaban -p azkaban_db < backup.sql从备份重建数据库先停止 Azkaban 服务
表结构升级使用 Azkaban 提供的 migration SQL 脚本升级时自动或手动执行存放于 sql/ 目录
注意事项- 备份文件应异地存储(如 S3)
- 定期验证备份可恢复性
- 执行前确保数据库连接池已关闭

10.5 监控指标采集(JMX / Prometheus Exporter)

指标类型JMX MBean 路径Prometheus Exporter 支持采集建议
JVM 基础java.lang:type=Memory
java.lang:type=Threading
支持监控堆内存、GC 次数、线程数
Web Server 状态azkaban.jmx:type=Server
activeExecutors, queueSize
需自定义 exporter关注任务队列积压情况
Executor 负载azkaban.exec.jmx:type=Executor
busyCount, idleCount
需开发实时反映资源利用率
HTTP 请求jetty:name=statistics,type=StatisticsHandler支持请求延迟、错误率
数据库连接com.mchange.v2.c3p0:type=PooledDataSource支持连接池使用率、等待数
自定义指标实现 MBean 注册通过 JMX Exporter 转换如失败任务数、调度延迟
采集工具- JConsole / VisualVM(调试)
- Prometheus + Grafana(生产)
使用 jmx_exporter.jar配置 jmx_exporter_config.yaml
注意事项- 开启 JMX 需配置 -Dcom.sun.management.jmxremote
- 生产环境建议启用认证与 SSL
- 高频采集可能影响性能

第十一章:REST API 编程接口

11.1 认证接口:获取 Session

方法请求方式参数示例注意事项
获取 Session IDPOST /
?action=login
&username=admin
&password=admin
用户名和密码curl -k -X POST "https://azkaban:8443/" -d "action=login&username=admin&password=admin"返回 JSON 包含 session.id
响应格式application/json-{"status": "success", "session.id": "abc123xyz"}失败返回 401 Unauthorized
Session 有效期默认 24 小时可通过 session.time.to.live.mins 配置超时后需重新登录建议程序化调用前先检查有效性
使用 Cookie 模式登录后保存 JSESSIONID用于后续请求自动认证curl -k -X POST ... -c cookie.txt推荐用于脚本批量操作
安全建议HTTPS 必须启用防止凭证泄露不在日志中打印密码生产环境禁用默认账户
注意事项- 所有 API 调用需携带 session.id 或 Cookie
- 多个 Web Server 共享同一 session 存储(数据库)
- 若启用 LDAP,此处使用 LDAP 凭证

11.2 项目操作接口:上传、删除、列表

操作请求方式参数示例注意事项
列出所有项目GET /manager?
ajax=fetchProjects
&session.id=abc123
session.idcurl -k "https://azkaban:8443/manager?ajax=fetchProjects&session.id=abc123"返回项目名、权限、描述
上传项目 ZIP 包POST /manager?
ajax=upload
&project=my_project
&session.id=abc123
project: 项目名
file: ZIP 文件
curl -k -X POST "https://azkaban:8443/manager?ajax=upload&project=etl_v2&session.id=abc123" -F "file=@etl_v2.zip"ZIP 必须包含 .job.flow 文件
删除项目GET /manager?
ajax=removeProject
&project=my_project
&session.id=abc123
project: 项目名curl -k "https://azkaban:8443/manager?ajax=removeProject&project=old_etl&session.id=abc123"永久删除,不可恢复
获取项目信息GET /manager?
ajax=fetchProjectInfo
&name=project_name
&session.id=abc123
name: 项目名查看创建者、权限等元数据用于审计或依赖分析
注意事项- 上传时若项目已存在且无 Write 权限则失败
- 删除项目需 Admin 权限
- ZIP 包大小受 max.request.size.mb 限制(默认 100MB)

11.3 执行控制接口:执行工作流、停止执行

操作请求方式参数示例注意事项
执行工作流POST /executor?
ajax=executeFlow
&project=etl
&flow=daily
&session.id=abc123
project, flowcurl -k -X POST "https://azkaban:8443/executor?ajax=executeFlow&project=etl&flow=daily&session.id=abc123"可附加 &param1=value1 传参
停止执行GET /executor?
ajax=cancelFlow
&execid=12345
&session.id=abc123
execid: 执行 IDcurl -k "https://azkaban:8443/executor?ajax=cancelFlow&execid=12345&session.id=abc123"向 Executor 发送中断信号
重新运行失败节点POST /executor?
ajax=executeFlow
&project=p1
&flow=f1
&execid=12345
&rerunFailed=true
rerunFailed=true仅重试失败分支原 execid 不变
强制重试指定节点POST /executor?
ajax=executeFlow
&project=p1
&flow=f1
&execid=12345
&nodes=nodeA,nodeB
nodes: 节点列表精确控制重试范围节点必须属于该流程
注意事项- 执行成功返回 execid
- 停止操作不保证立即生效(取决于任务是否可中断)
- 无法停止已完成的 execid

11.4 状态查询接口:获取执行详情、日志片段

操作请求方式参数示例注意事项
获取执行状态GET /executor?
ajax=fetchExecStatus
&execid=12345
&session.id=abc123
execidcurl -k "https://azkaban:8443/executor?ajax=fetchExecStatus&execid=12345&session.id=abc123"返回整体状态(RUNNING/SUCCEEDED/FAILED)
获取节点状态详情GET /executor?
ajax=fetchExecFlow
&execid=12345
&session.id=abc123
execid包含每个 Job 的开始/结束时间、主机、返回码用于 DAG 渲染
获取日志片段GET /executor?
ajax=fetchExecJobLogs
&execid=12345
&jobId=jobA
&offset=0
&length=10000
&session.id=abc123
offset, length 控制分页offset=0&length=10000 获取前 10KB大日志需多次请求
获取所有日志(完整)循环调用 fetchExecJobLogs直到 offset >= total脚本中拼接所有片段注意性能影响
获取最近 N 次执行GET /executor?
ajax=fetchFlowExecutions
&project=etl
&flow=daily
&start=0
&length=5
start, length 分页查询历史执行结果用于趋势分析
注意事项- 日志接口返回纯文本或 JSON
- offset 超出范围返回空
- 建议设置超时和重试机制

11.5 调度管理接口:创建、更新、删除定时任务

操作请求方式参数示例注意事项
创建调度POST /schedule?
ajax=scheduleFlow
&project=etl
&flow=daily
&cronExpression=0+0+2+++%3F
&session.id=abc123
cronExpression URL 编码%3F 表示 ?可附加 &param1=value1
获取调度详情GET /schedule?
ajax=fetchSchedule
&scheduleId=67890
&session.id=abc123
scheduleId 来自创建响应返回下一次触发时间用于验证调度是否生效
更新调度POST /schedule?
ajax=reschedule
&scheduleId=67890
&cronExpression=0+30+2+++%3F
&session.id=abc123
scheduleId + 新表达式修改执行时间原参数保留
删除调度GET /schedule?
ajax=unschedule
&scheduleId=67890
&session.id=abc123
scheduleId彻底移除定时任务不影响手动执行
列出项目所有调度GET /schedule?
ajax=fetchSchedulesByProject
&project=etl
&session.id=abc123
project返回多个 scheduleId用于批量管理
注意事项- cronExpression 需 URL 编码(空格→+?%3F
- 创建调度需项目 Write 权限
- 调度状态存储于 MySQL,重启不影响

第十二章:最佳实践与常见问题

12.1 工作流命名规范与模块化设计

实践建议示例说明
命名规范使用小写字母、下划线、数字daily_user_sync, hourly_metrics避免空格和特殊字符
层级命名按业务域+频率+功能划分finance_monthly_close, log_hourly_cleanup便于搜索和分类
模块化设计将通用流程抽为子项目utils_project → vacuum_db_flow通过 embeddedFlow 复用
版本控制在项目名或描述中标注版本etl_pipeline_v2, v2025q3配合 Git 管理变更
文档化project.properties 中添加 descriptiondescription=Daily ETL for CRM data提高可维护性
注意事项- 避免单个工作流超过 50 个节点
- 使用嵌套流程实现分层调度
- 定期重构冗余 Job

12.2 异常处理与告警集成(邮件、Webhook)

方法配置/实现方式示例注意事项
邮件告警配置 azkaban.properties
mail.sender=admin@company.com
mail.host=smtp.company.com
mail.user=azkaban
mail.password=***
在失败 Job 后添加 send_email.job需启用 SMTP 支持
Webhook 通知使用 command job 调用 curlcommand=curl -X POST https://hooks.slack.com/services/... -d '{"text":"Flow failed!"}'支持 Slack、钉钉、企业微信
失败回调机制在 Job 中捕获 exit code 并触发通知if [ $? -ne 0 ]; then send_alert; exit 1; fi实现细粒度控制
全局失败处理设置 on.failure=END 并添加通知节点所有失败路径汇聚到 notify_fail Job统一告警入口
静默期设置在通知脚本中判断时间如夜间 00:00–06:00 不发送短信避免打扰
注意事项- 敏感信息(如 webhook token)使用 secure properties
- 告警应包含 execid、flow 名、失败节点
- 避免告警风暴(如重复失败)

12.3 版本升级注意事项

阶段操作建议注意事项
升级前准备- 备份数据库(MySQL)
- 备份 executor logs 和 conf 目录
- 查看 release notes 兼容性
确认新版本是否支持当前 Hadoop 版本
停止服务先停止 Executor,再停止 Web Server避免执行中任务中断
数据库迁移运行 Azkaban 提供的 SQL 升级脚本sql/upgrade-to-3.0.0.sql
配置文件更新对比新旧 azkaban.properties新增参数需手动添加
插件兼容性检查自定义 job type 是否兼容重新编译或替换 jar 包
升级后验证- 启动服务
- 登录 UI
- 执行测试 flow
- 检查日志无 ERROR
观察 JMX 指标是否正常
回滚计划准备旧版本包和备份若失败立即回滚
注意事项- 建议在维护窗口升级
- 先在测试环境验证
- 升级后清除浏览器缓存

12.4 大规模作业调度优化建议

优化方向建议措施说明
减少 DAG 复杂度拆分大流程为多个子流程提高并行度和可维护性
合理设置并发控制同时运行的 exec 数避免数据库或 Hadoop 资源过载
使用 Executor 分组按业务线或优先级分组关键任务独占资源
优化依赖结构消除不必要的依赖允许更多任务并行
日志归档策略定期清理旧执行日志防止磁盘溢出
数据库优化executionsexecution_jobs 表建立索引加快查询速度
异步上传使用 API 批量上传项目减少 UI 操作延迟
注意事项- 避免”巨型工作流”(>100 节点)
- 监控 queueSize 和 executor busyCount
- 使用 embeddedFlow 替代 flat DAG

12.5 典型故障场景与解决方案

故障现象可能原因解决方案预防措施
Web UI 无法登录密码错误、LDAP 连接失败、session 过期检查认证配置,重启服务启用多因素认证
工作流卡在 RUNNINGExecutor 断开、任务无响应查看 Executor 状态,手动 cancel 后重试设置 job 超时(ulimit)
日志无法查看Executor 磁盘满、网络不通清理日志目录,检查防火墙配置日志轮转
调度未触发Cron 表达式错误、Web Server 故障检查 schedule 表,重启 Web Server使用 UI 生成 Cron
上传项目失败ZIP 包损坏、权限不足、大小超限重新打包,检查 Write 权限验证 ZIP 并压缩
数据库连接池耗尽连接未释放、并发过高增加 maxPoolSize,优化查询监控连接数,设置超时
Executor 注册失败网络不通、端口冲突、配置错误检查 executor.portweb.server.url使用健康检查脚本
注意事项- 建立故障知识库(Wiki)
- 所有变更记录操作日志
- 关键节点添加健康检查 Job