Article
第一章: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)
| 对比维度 | Azkaban | Apache Airflow | Apache 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=commandcommand=echo "Hello"description=Print hello message | - 每行一个 key=value- 不支持嵌套结构 |
| 依赖定义语法 | jobB:type=commandcommand=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 | 展开目录结构 | (略) | 目录包含 bin、lib、plugins 等子目录 |
| 启动服务 | bin/start-solo.sh | 启动内置 Web + H2 + Executor | (无参数) | 自动打开 http://localhost:8081 |
| 停止服务 | bin/shutdown-solo.sh | 安全关闭服务 | (无参数) | 避免直接 kill 进程 |
| 访问 Web UI | 浏览器访问 http://<host>:8081 | 登录管理界面 | 默认账号:azkaban / azkaban | H2 数据存储在本地,重启可能丢失 |
2.2 双服务模式(Two-Server Mode)部署
| 方法/步骤 | 语法/命令 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| 准备 Web Server | 解压 azkaban-web-server-* | 部署 Web 服务 | (略) | 需配置指向 MySQL 和 Executor 地址 |
| 准备 Executor Server | 解压 azkaban-executor-server-* | 部署执行服务 | (略) | 每个 Executor 需独立端口(默认 12321) |
| 配置 database.properties | driver=com.mysql.jdbc.Driveruser=azkabanpassword=azkabanurl=jdbc:mysql://localhost:3306/azkaban | Web 和 Executor 共用数据库配置 | 放置于 conf/ 目录下 | 必须提前创建数据库和用户 |
| 初始化数据库 | mysql -u root -p < scripts/create-all-sql-3.84.4.sql | 创建表结构 | 替换脚本版本号 | 脚本位于 azkaban-db 模块 |
| 启动 Web Server | bin/start-web.sh | 启动 Web 服务 | (无参数) | 依赖 Executor 已注册 |
| 启动 Executor Server | bin/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.properties | driver=com.mysql.jdbc.Driveruser=azkabanpassword=azkabanurl=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 -tlnp | grep 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-8 | 在 start-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_dailyDescription: 每日ETL任务 | 项目名不可重复,不支持特殊字符 |
| 填写项目属性 | 设置项目描述、负责人、权限组 | 便于团队协作与审计 | Owner: data_team | 描述建议清晰说明用途 |
| 提交创建 | 点击 “Create” 按钮 | 在数据库中持久化项目元数据 | (无返回值) | 成功后跳转至项目主页 |
3.2 使用命令行工具(azkaban-cli)上传项目
| 方法 | 语法/命令 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| 打包项目目录 | zip -r project.zip *.job *.properties | 将 job 文件打包为 zip | 包含所有依赖脚本和配置文件 | 必须为 zip 格式,不支持 tar/gz |
| 获取 Session ID | curl -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_name 和 session.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 文件基础语法结构
| 属性 | 语法格式 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| type | type=command | java | hive | flow 等 | 定义任务类型 | type=command | 必须指定,决定执行器行为 |
| command | command=echo "hello" | 执行 Shell 命令 | command=hive -f /path/to/script.hql | 多行命令可用 \ 换行 |
| description | description=This is a test job | 添加注释说明 | 提高可读性 | Web UI 中显示 |
| retries | retries=3 | 失败自动重试次数 | 每次间隔由 retry.backoff 值决定 | 不适用于关键一致性任务 |
| retry.backoff | retry.backoff=10000 | 重试间隔(毫秒) | 默认 10 秒 | 避免频繁重试压垮系统 |
| memory / cpu | memory.mb=2048cpu.threads=2 | 资源限制(需插件支持) | 控制作业资源占用 | 非强制,依赖外部监控 |
| 注释 | # 这是一行注释 | 忽略该行内容 | 支持行尾注释 | 不支持多行 /* */ |
4.2 定义节点依赖关系(dependencies)
| 方法 | 语法 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| 单依赖 | dependencies=jobA | 当 jobA 成功后执行当前任务 | jobB:type=commandcommand=echo "B"dependencies=jobA | 依赖任务必须在同一项目中 |
| 多依赖 | dependencies=jobA,jobC | 所有前置任务成功后才运行 | jobD:type=commandcommand=merge_data.shdependencies=extract1,extract2 | 逗号分隔,无空格 |
| 并行执行 | 多个任务无依赖关系 | 同时启动多个分支 | jobA → jobCjobB → jobC | 提升整体执行效率 |
| 循环依赖检测 | 系统自动检测 | 防止死锁 | jobA → jobB → jobA ❌ | 提交时报错 “Circular dependency” |
| 依赖命名规范 | 使用有意义的名称 | 提高可维护性 | extract_user_data, transform_sales | 避免使用 job1, job2 |
4.3 使用内联工作流(Inline Flow in job file)
| 方法 | 语法 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| 定义内联 Flow | type=flowflow.name=myFlownodes=[{...}] | 在 job 文件中直接定义复杂 DAG | 见下方示例 | 适合小型嵌套流程 |
| nodes 数组结构 | nodes=[{"name":"A","type":"command","command":"echo A"}, {"name":"B","dependencies":["A"],"command":"echo B"}] | 描述 DAG 节点 | JSON 格式嵌入 | 注意引号转义 |
| 嵌套 Flow 调用 | type=embeddedFlowflow.project=project_xflow.flowId=flow_y | 调用其他项目的 Flow | 实现模块复用 | 调用方需有 Read 权限 |
| 示例:简单内联 Flow | myflow.job:type=flowflow.name=inline_testnodes=[{"name":"start","type":"command","command":"echo start"},{"name":"end","dependencies":["start"],"type":"command","command":"echo done"}] | 替代多个 .job 文件 | 减少文件数量 | 复杂逻辑仍建议拆分 |
| 注意事项 | - JSON 必须合法,避免语法错误 - 不支持跨项目内联定义 |
4.4 复杂 DAG 图构建技巧
| 技巧 | 说明 | 示例 | 注意事项 |
|---|---|---|---|
| 分层设计 | 按阶段划分:Extract → Transform → Load | 每层多个并行任务汇聚到下一阶段 | 降低耦合度,便于调试 |
| 子流程复用 | 将通用逻辑封装为独立 Flow | 如 daily_cleanup.flow | 通过 embeddedFlow 调用 |
| 动态参数传递 | 使用 ${input_path} 在运行时替换 | 配合执行时传参实现灵活调度 | 参数需提前声明或传入 |
| 错峰执行 | 设置高耗时任务错开时间 | 避免资源争抢 | 使用调度延迟或优先级 |
| 虚拟节点(Dummy Job) | 添加空任务作为汇合点 | wait_all:type=noopdependencies=task1,task2 | noop 类型不执行实际操作 |
4.5 错误处理与失败策略(fail-action, on-failure)
| 属性 | 语法 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| fail-action | fail-action=finishCurrent | 定义失败后的整体行为 | finishCurrent | cancelImmediately | retry | 作用于整个 Flow |
| on-failure | on.failure=finishCurrent | 同 fail-action,旧写法 | 推荐统一使用 fail-action | 两者功能相同,优先用新语法 |
| finishCurrent | fail-action=finishCurrent | 当前任务失败后,允许已完成任务继续,不启动新任务 | 适用于非关键路径任务 | 最终 Flow 状态为 FAILED |
| cancelImmediately | fail-action=cancelImmediately | 一旦失败立即终止所有运行和待运行任务 | 关键任务链使用 | 防止脏数据传播 |
| retry(节点级) | retries=3retry.backoff=5000 | 单个任务自动重试 | 适用于瞬时故障(如网络抖动) | 不适用于数据一致性破坏场景 |
| 自定义失败处理 Job | job_failure_handler:type=commandcommand=notify_error.sh ${azkaban.flow.execid}on.failure=END | 失败时执行通知脚本 | 发送邮件或 Webhook | 需确保通知服务可靠 |
第五章:Job 类型与处理器(Job Types)
5.1 Command Job Type(command)
| 方法/属性 | 语法 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| type | type=command | 定义为 Shell 命令任务 | type=command | 必须指定 |
| command | command=<shell_command> | 执行任意 Shell 命令 | command=echo "Hello"command=sh /data/scripts/cleanup.sh | 支持多行(用 \ 连接) |
| working.dir | working.dir=/path/to/dir | 设置工作目录 | working.dir=/home/azkaban/jobs | 默认为执行器临时目录 |
| env. | env.PATH=/usr/local/bin:$PATH | 设置环境变量 | env.HADOOP_HOME=/opt/hadoop | 仅对该 Job 有效 |
| ulimit | ulimit=-v 8388608 | 限制内存使用(KB) | ulimit=-u 64 | 防止资源耗尽 |
| 注意事项 | - 不支持交互式命令 - 建议将脚本外部化而非写在 job 文件中 |
5.2 Java Job Type(java)
| 方法/属性 | 语法 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| type | type=java | 调用本地 JVM 执行 Java 类 | type=java | 类必须可被 classpath 加载 |
| class | class=com.example.MainClass | 指定主类名 | class=com.etl.BatchProcessor | 必须包含 main 方法 |
| jars | jars=/path/to/lib/a.jar,/path/to/lib/b.jar | 添加依赖 JAR 包路径 | jars=/home/azkaban/lib/utils.jar | 多个用逗号分隔 |
| args | args=arg1 arg2 | 传入 main 方法参数 | args=input.parquet output.orc | 空格分隔多个参数 |
| jvm.args | jvm.args=-Xmx2g -Dlog.level=INFO | 设置 JVM 启动参数 | jvm.args=-server -XX:+UseG1GC | 影响性能和稳定性 |
| lib.dir | lib.dir=/path/to/libs | 自动加载该目录下所有 jar | lib.dir=/home/azkaban/java_libs | 简化依赖管理 |
| 注意事项 | - 主类需打包进 jar 并确保无冲突依赖 - 推荐使用 hadoopJava 处理 Hadoop 作业 |
5.3 Pig Job Type(pig)
| 方法/属性 | 语法 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| type | type=pig | 使用本地 Pig 引擎运行脚本 | type=pig | 已废弃,建议用 hadoopPig |
| pig.script | pig.script=/path/to/script.pig | 指定 Pig 脚本路径 | pig.script=/home/user/analyze.pig | 脚本必须存在且可读 |
| parameters | parameters="input=data.csv output=result" | 传入参数 | parameters="year=2025 debug=true" | 双引号包裹多个参数 |
| 注意事项 | - 该类型不连接 Hadoop 集群 - 仅用于测试或单机模式 |
5.4 Hadoop Pig Job Type(hadoopPig)
| 方法/属性 | 语法 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| type | type=hadoopPig | 在 Hadoop 集群上运行 Pig 脚本 | type=hadoopPig | 需配置 Hadoop 环境 |
| pig.script | pig.script=/path/to/script.pig | 指定 Pig 脚本 | pig.script=hdfs://nn:9000/pig/etl.pig | 支持 HDFS 路径 |
| hadoop.security.manager.class | hadoop.security.manager.class=org.apache.azkaban.jobtype.HadoopSecurityManager_H_2_0 | 指定安全管理器 | 根据 Hadoop 版本选择 | Azkaban 插件需支持 |
| job.flow.id | job.flow.id=${azkaban.flow.execid} | 传递 Flow ID(可选) | 用于日志追踪 | 通常自动注入 |
| parameters | parameters="in=hdfs://... out=hdfs://..." | 传入 Pig 参数 | parameters="date=${today}" | 支持变量替换 |
| 注意事项 | - Executor 必须安装 Pig 并配置好 Hadoop 客户端 - 脚本中可用 $parameter 接收参数 |
5.5 Hadoop WordCount Job Type(hadoopWordCount)
| 方法/属性 | 语法 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| type | type=hadoopWordCount | 简化版 MapReduce 任务模板 | type=hadoopWordCount | 专为 WordCount 示例设计 |
| input.path | input.path=/data/input.txt | 输入文件路径 | input.path=hdfs://nn:9000/input | 支持通配符如 *.txt |
| output.path | output.path=/data/output | 输出目录(必须不存在) | output.path=hdfs://nn:9000/out/2025 | 任务失败前会自动删除 |
| hadoop.job.class | hadoop.job.class=com.hadoop.mapreduce.WordCount | 自定义 MapReduce 类 | 替换默认实现 | 必须继承相应接口 |
| 注意事项 | - 实际生产中不推荐使用此类型 - 更通用的做法是使用 hadoopJava 或直接提交 Jar |
5.6 Hive Job Type(hive)
| 方法/属性 | 语法 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| type | type=hive | 使用本地 Hive CLI 执行 HQL | type=hive | 不连接 HiveServer2,已过时 |
| hive.script | hive.script=/path/to/query.hql | 指定 HQL 脚本路径 | hive.script=/home/azkaban/sql/dim_user.hql | 文件需 UTF-8 编码 |
| hive.home | hive.home=/opt/hive | 设置 Hive 安装路径 | hive.home=/usr/lib/hive | 影响 bin/hive 调用 |
| hive.parameters | hive.parameters="dt=2025-01-01 env=prod" | 传入 Hive 变量 | SELECT * FROM logs WHERE dt='${dt}' | 在 HQL 中用 ${} 引用 |
| 注意事项 | - 建议改用 hadoopHive 类型以支持远程执行- 本地模式需安装 Hive 客户端 |
5.7 EmbeddedFlow Job Type(嵌套流程)
| 方法/属性 | 语法 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| type | type=embeddedFlow | 调用当前项目或其他项目的子流程 | type=embeddedFlow | 实现模块化设计 |
| flow.project | flow.project=my_subproject | 指定目标项目名 | flow.project=common_libs | 省略则表示当前项目 |
| flow.flowId | flow.flowId=cleanup_flow | 指定要调用的工作流 ID | flow.flowId=daily_export | 必须存在且可访问 |
| propagate.failure | propagate.failure=true | 是否传播子流程失败 | true | false | 设为 true 则父流程也失败 |
| flow.parameters | flow.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.properties | name=myCustomTypeclass=com.example.MyJobRunner | 注册新 Job 类型 | name 是 job 文件中使用的 type 值 |
| 实现 JobRunner | public 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} | 当前工作流的 ID | daily_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.parameters | flow.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.file | secure.props.file=/path/to/secure.properties | 指定加密属性文件路径 | 放置于 executor conf/ 目录 | Web Server 不需要 |
| secure.property.names | secure.property.names=db.password,api.key | 声明哪些参数为敏感信息 | 多个用逗号分隔 | 必须匹配实际参数名 |
| 传参方式 | 在执行时传入 db.password=xxx | 但日志和 UI 中隐藏 | 显示为 ****** | 防止密码泄露 |
| 属性文件格式 | db.password=encrypted(AES):base64dataapi.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-30dry_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=prodbatch_size=10000 | 参数对所有调度实例生效 |
| 保存调度 | 点击 “Schedule” 按钮 | 提交调度任务到数据库 | 成功后显示下一次执行时间 | 调度信息存储在 MySQL 中 |
| 注意事项 | - 调度由 Web Server 统一管理 - Executor 不参与调度决策 - 修改调度需重新提交 |
7.3 Cron 表达式在 Azkaban 中的应用
| 字段位置 | 含义 | 允许值 | 示例 | 说明 |
|---|---|---|---|---|
| 1 | 秒 | 0–59 | 0 | Azkaban 默认从分钟开始,秒固定为 0 |
| 2 | 分钟 | 0–59 | 0 | 每小时整点触发 |
| 3 | 小时 | 0–23 | 2 | 凌晨 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 或消息队列通知 Azkaban | Kafka 消息 → 调用 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 的 Duration | Execution 页面时间轴 | 找出最长耗时任务,优化 SQL 或脚本 |
| 并行度不足 | DAG 中存在长链无分支 | 优化依赖结构 | 拆分独立任务并行执行 |
| 资源竞争 | 多个 Job 同时访问同一数据库或 HDFS 路径 | 系统监控工具(如 Grafana) | 错峰调度或增加资源 |
| Executor 负载过高 | 某个 Executor 长时间 busy | Web 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=trueldap.hostname=ldap.company.comldap.port=389ldap.base.dn="ou=users,dc=company,dc=com" | 集中管理企业用户 | 支持 Active Directory | 需网络可达且证书可信 |
| JDBC 认证 | 配置数据库连接信息:database.type=mysqldatabase.host=localhostdatabase.user=azkabandatabase.password=***并实现自定义 UserDAO | 与现有系统集成 | 用于已有用户中心的场景 | 需开发插件支持 |
| 启用认证类 | azkaban.user.manager=azkaban.user.XmlUserManager | 指定用户管理器实现 | File: XmlUserManagerLDAP: 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 支持
| 配置项 | 语法/路径 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| 启用 HTTPS | 在 azkaban.properties 中设置:jetty.ssl.context.key.store.path=/path/to/keystore.jksjetty.ssl.context.key.store.password=changeitjetty.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=8443 | HTTPS 监听端口 | 默认关闭 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.propertiessecure.property.names=db.password,api.key | 声明敏感参数名 | 执行时传入但日志中隐藏 | 仅 Executor 能解密 |
| 加密属性文件 | db.password=encrypted(AES):base64dataapi.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=AteamB_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=8executor.resources.mem=16g | 需自定义调度器 |
| Executor 分组 | 使用 executor.group=etl 或 realtime | 按任务类型隔离 | 将批处理与实时任务分开 | 提高稳定性 |
| 查看负载状态 | 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)
| 组件 | 配置要点 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| HAProxy | frontend azkaban_httpsbind *:8443 ssl crt /certs/azkaban.pemdefault_backend azkaban_serversbackend azkaban_serversbalance roundrobinserver web1 192.168.1.10:8443 check ssl verify noneserver web2 192.168.1.11:8443 check ssl verify none | TCP 层负载均衡 | 支持 SSL 终止或透传 | verify none 用于自签名证书 |
| Nginx | upstream 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-binserver-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=Memoryjava.lang:type=Threading | 支持 | 监控堆内存、GC 次数、线程数 |
| Web Server 状态 | azkaban.jmx:type=ServeractiveExecutors, queueSize | 需自定义 exporter | 关注任务队列积压情况 |
| Executor 负载 | azkaban.exec.jmx:type=ExecutorbusyCount, 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 ID | POST /?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.id | curl -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, flow | curl -k -X POST "https://azkaban:8443/executor?ajax=executeFlow&project=etl&flow=daily&session.id=abc123" | 可附加 ¶m1=value1 传参 |
| 停止执行 | GET /executor?ajax=cancelFlow&execid=12345&session.id=abc123 | execid: 执行 ID | curl -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 | execid | curl -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 表示 ? | 可附加 ¶m1=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 中添加 description | description=Daily ETL for CRM data | 提高可维护性 |
| 注意事项 | - 避免单个工作流超过 50 个节点 - 使用嵌套流程实现分层调度 - 定期重构冗余 Job |
12.2 异常处理与告警集成(邮件、Webhook)
| 方法 | 配置/实现方式 | 示例 | 注意事项 |
|---|---|---|---|
| 邮件告警 | 配置 azkaban.properties:mail.sender=admin@company.commail.host=smtp.company.commail.user=azkabanmail.password=*** | 在失败 Job 后添加 send_email.job | 需启用 SMTP 支持 |
| Webhook 通知 | 使用 command job 调用 curl | command=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 分组 | 按业务线或优先级分组 | 关键任务独占资源 |
| 优化依赖结构 | 消除不必要的依赖 | 允许更多任务并行 |
| 日志归档策略 | 定期清理旧执行日志 | 防止磁盘溢出 |
| 数据库优化 | 对 executions、execution_jobs 表建立索引 | 加快查询速度 |
| 异步上传 | 使用 API 批量上传项目 | 减少 UI 操作延迟 |
| 注意事项 | - 避免”巨型工作流”(>100 节点) - 监控 queueSize 和 executor busyCount - 使用 embeddedFlow 替代 flat DAG |
12.5 典型故障场景与解决方案
| 故障现象 | 可能原因 | 解决方案 | 预防措施 |
|---|---|---|---|
| Web UI 无法登录 | 密码错误、LDAP 连接失败、session 过期 | 检查认证配置,重启服务 | 启用多因素认证 |
| 工作流卡在 RUNNING | Executor 断开、任务无响应 | 查看 Executor 状态,手动 cancel 后重试 | 设置 job 超时(ulimit) |
| 日志无法查看 | Executor 磁盘满、网络不通 | 清理日志目录,检查防火墙 | 配置日志轮转 |
| 调度未触发 | Cron 表达式错误、Web Server 故障 | 检查 schedule 表,重启 Web Server | 使用 UI 生成 Cron |
| 上传项目失败 | ZIP 包损坏、权限不足、大小超限 | 重新打包,检查 Write 权限 | 验证 ZIP 并压缩 |
| 数据库连接池耗尽 | 连接未释放、并发过高 | 增加 maxPoolSize,优化查询 | 监控连接数,设置超时 |
| Executor 注册失败 | 网络不通、端口冲突、配置错误 | 检查 executor.port 和 web.server.url | 使用健康检查脚本 |
| 注意事项 | - 建立故障知识库(Wiki) - 所有变更记录操作日志 - 关键节点添加健康检查 Job |