ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Airflow 和 DolphinScheduler 工作流讲解

Airflow 和 DolphinScheduler 工作流讲解 Airflow 和 Apache DolphinScheduler 都是数据工程中常用的工作流调度平台核心目标都是按照预先定义的依赖关系在指定时间或事件触发任务并负责执行、重试、监控和告警。二者都可以用于离线数仓、ETL、机器学习流程、数据同步和批处理任务但设计理念和使用方式有所不同。一、Airflow 工作流1. 核心概念Airflow 的工作流称为DAG即 Directed Acyclic Graph有向无环图。一个 DAG 由多个 Task 组成Task 之间通过依赖关系连接数据采集 ↓ 数据清洗 ↓ 数据入库 ↓ 指标计算 ↓ 报表生成Airflow 通常使用 Python 代码定义工作流fromairflowimportDAGfromairflow.operators.bashimportBashOperatorfromdatetimeimportdatetimewithDAG(dag_iddaily_etl,start_datedatetime(2024,1,1),schedule0 2 * * *,catchupFalse,)asdag:extractBashOperator(task_idextract,bash_commandpython extract.py,)transformBashOperator(task_idtransform,bash_commandpython transform.py,)loadBashOperator(task_idload,bash_commandpython load.py,)extracttransformload这里定义了一个每天凌晨 2 点执行的 ETL 流程extract → transform → load2. Airflow 的执行流程Airflow 的典型执行过程如下调度器 Scheduler ↓ 发现需要执行的 DAG ↓ 检查任务依赖 ↓ 提交任务到执行器 Executor ↓ Worker 执行任务 ↓ 记录任务状态和日志 ↓ 成功、失败、重试或触发下游任务主要组件包括组件作用Scheduler负责判断哪些 DAG 和任务应该执行Webserver提供 Web 管理界面Metadata Database保存 DAG、任务状态、日志索引等元数据Executor决定任务如何被执行Worker实际执行任务DAG Processor解析 Python DAG 文件常见 ExecutorSequentialExecutor单进程执行适合测试LocalExecutor本机并行执行CeleryExecutor通过 Celery 分布式执行KubernetesExecutor每个任务使用 Kubernetes Pod 执行3. Airflow 的任务类型Airflow 使用 Operator 表示任务。常见 OperatorBashOperator 执行 Shell 命令 PythonOperator 执行 Python 函数 SQL Operator 执行 SQL DockerOperator 运行 Docker 容器 KubernetesPodOperator 运行 Kubernetes Pod HttpOperator 调用 HTTP 接口 Sensor 等待某个条件满足例如等待文件产生文件 Sensor ↓ 读取文件 ↓ 处理数据Airflow 也有大量 Provider支持MySQL、PostgreSQL、HiveSpark、FlinkHadoop、YARNKafkaAWS、Azure、GCPDocker、Kubernetesdbt 等4. Airflow 的特点优点Python 原生灵活性很高生态成熟扩展丰富适合复杂的数据工程和机器学习流程对任务逻辑、动态任务、条件分支支持较好可以通过代码进行版本管理和测试缺点学习成本相对较高部署和运维复杂度较高很多功能需要通过 Python 编码实现对大规模任务的调度和元数据数据库压力需要重点优化Airflow 更偏向“工作流编排平台”不是传统意义上的全功能数据开发平台二、DolphinScheduler 工作流DolphinScheduler 中文通常称为“海豚调度”它的核心对象是工作流定义也常称为 Process Definition。DolphinScheduler 更偏向于通过 Web 页面配置工作流也支持通过代码或 API 创建工作流。1. 工作流示例Shell 任务数据采集 ↓ SQL 任务数据清洗 ↓ Spark 任务计算指标 ↓ 通知任务发送结果在 DolphinScheduler 中用户可以拖拽任务节点并连接任务之间的依赖关系。常见任务类型ShellSQLPythonSparkFlinkMapReduceHiveDataXHTTPSub-ProcessSwitchDependentConditionsNotifications2. DolphinScheduler 的执行流程Master Server ↓ 接收工作流实例 ↓ 拆分并提交任务 ↓ Worker Group 执行任务 ↓ 记录日志、状态和结果 ↓ 继续执行满足条件的后续任务主要组件包括组件作用API Server对外提供工作流管理接口Master Server负责工作流调度、依赖判断和任务分发Worker Server执行具体任务Alert Server处理失败告警和通知Database保存工作流和任务元数据Registry Center保存服务注册和节点信息常使用 ZooKeeperDolphinScheduler 一般会按照 Worker Group 对任务进行分组例如default 普通任务 spark-workers Spark 任务 flink-workers Flink 任务 gpu-workers 机器学习任务这样可以将不同类型的任务分配到不同资源池中。3. DolphinScheduler 的典型配置流程通常步骤如下创建租户创建用户和项目配置数据源创建工作流添加任务节点配置任务参数连接任务依赖设置定时规则发布工作流启动或定时运行例如数据源配置 ↓ 创建项目 ↓ 创建工作流 ↓ 添加 Shell、SQL、Spark 节点 ↓ 配置依赖关系 ↓ 发布工作流 ↓ 设置定时上线4. DolphinScheduler 的特点优点Web 拖拽式开发上手较快对传统大数据组件支持比较直接适合企业数据平台和数仓调度自带租户、项目、用户、资源、告警等管理能力Worker Group 方便做资源隔离对 Shell、SQL、Spark、Flink、DataX 等任务支持方便工作流发布、上线、下线流程较清晰缺点复杂任务逻辑的灵活性通常不如 Python 化的 Airflow复杂动态工作流编排时页面配置可能变得繁琐生态和第三方 Operator 数量相对 Airflow 少大量工作流通过页面维护时代码版本管理能力需要额外建设复杂部署同样需要考虑 ZooKeeper、数据库、Master、Worker 等组件三、两者的核心区别对比项AirflowDolphinScheduler工作流定义主要通过 Python 代码主要通过 Web 页面也支持 API核心对象DAGProcess Definition使用体验偏开发者偏数据平台和调度运维人员灵活性很高较高可视化配置有界面但主要用于查看和管理拖拽式配置较强任务扩展Provider、Operator 生态丰富内置大数据任务类型丰富资源隔离Executor、Queue、Pool、KubernetesWorker Group、队列、优先级版本管理Python 文件天然适合 Git页面配置需要结合导出、API 或代码化方案条件分支Python 和 Branch Operator 灵活Switch、Conditions 等节点动态任务Dynamic Task Mapping 支持较好可以实现但方式相对不同大数据任务需要配置对应 Operator 或 ProviderSpark、Flink、Hive、DataX 等较直接学习成本较高页面入门较快典型场景复杂编排、数据工程、ML Pipeline企业数仓、大数据任务统一调度四、用一个实际例子对比假设每天需要执行1. 从业务库抽取订单数据 2. 将数据写入 HDFS 3. 使用 Spark 清洗数据 4. 使用 Hive 生成汇总表 5. 使用 DataX 同步到 MySQL 6. 任务失败时发送钉钉通知Airflow 实现思路MySQL ExtractOperator ↓ HDFS Upload Task ↓ SparkSubmitOperator ↓ HiveOperator ↓ DataX BashOperator ↓ Webhook Notification工作流依赖通过 Python 编写extractuploadsparkhivedataxnotify适合需要根据业务逻辑动态决定执行路径的场景例如如果数据量大: 使用 Spark 否则: 使用 PythonDolphinScheduler 实现思路在页面中添加SQL 节点 ↓ Shell 节点 ↓ Spark 节点 ↓ Hive 节点 ↓ DataX 节点 ↓ 告警节点然后设置每天执行时间失败重试次数超时时间Worker Group告警策略任务优先级资源参数对于标准化数仓流程这种方式比较直观。五、如何理解 Airflow 的 DAGAirflow 的 DAG 不只是“任务列表”而是一个完整的调度定义通常包括任务节点 任务依赖 调度周期 开始时间 重试策略 超时设置 并发限制 资源池 失败处理 告警配置例如withDAG(dag_idorder_etl,schedule0 1 * * *,catchupFalse,max_active_runs1,default_args{retries:2,retry_delay:timedelta(minutes5),},):...含义是每天凌晨 1 点执行不补跑历史周期同时最多运行一个 DAG 实例任务失败后最多重试 2 次每次重试间隔 5 分钟六、如何理解 DolphinScheduler 的工作流实例DolphinScheduler 中通常需要区分工作流定义 ↓ 定时调度 ↓ 工作流实例 ↓ 任务实例例如工作流定义订单每日处理 ↓ 2024-01-01 工作流实例 ↓ 2024-01-01 的抽取任务、清洗任务、汇总任务 ↓ 2024-01-02 工作流实例 ↓ 2024-01-02 的抽取任务、清洗任务、汇总任务这里工作流定义描述流程应该怎么执行工作流实例某一次实际运行任务实例该次运行中的具体任务执行记录这与 Airflow 中的DAG DAG Run Task Instance概念基本对应。七、如何选择适合选择 Airflow 的情况选择 Airflow 通常是因为团队 Python 能力较强工作流逻辑复杂需要动态生成任务需要大量调用 API、模型、数据处理代码已经使用 Kubernetes、云服务或现代数据栈希望所有调度逻辑通过 Git 管理需要丰富的第三方 Provider典型场景数据采集 API 调用 Python 处理 机器学习 模型部署适合选择 DolphinScheduler 的情况选择 DolphinScheduler 通常是因为主要面向企业数据仓库有较多 Shell、SQL、Hive、Spark、Flink、DataX 任务希望非开发人员也能配置流程需要完善的租户、项目、权限和告警管理希望使用页面快速搭建工作流需要对不同任务分配不同 Worker 资源组典型场景数据抽取 → ODS → DWD → DWS → ADS → 报表同步八、二者的共同运行机制无论使用 Airflow 还是 DolphinScheduler一个工作流任务通常都会经历这些状态等待调度 ↓ 准备执行 ↓ 运行中 ↓ 成功失败时可能进入运行失败 ↓ 自动重试 ↓ 重试成功 / 最终失败典型状态包括待运行运行中成功失败暂停跳过阻塞超时取消调度系统一般还需要处理上游任务是否成功是否满足定时条件是否达到并发限制是否有可用 Worker是否超过重试次数是否满足数据依赖是否发送告警九、简单总结可以用一句话概括Airflow 更像“用 Python 编写的工作流编排框架”DolphinScheduler 更像“面向企业数据平台的可视化工作流调度系统”。如果主要是复杂逻辑、Python、API、机器学习、动态编排可以优先考虑 Airflow。如果主要是数仓任务、SQL、Shell、Spark、Flink、DataX、权限和资源管理可以优先考虑 DolphinScheduler。二者也可以同时存在例如DolphinScheduler 负责企业级批处理调度 Airflow 负责复杂的数据科学和机器学习流程但需要避免两个系统重复调度同一批任务否则容易造成重复执行、依赖难以追踪和告警责任不清。
返回列表