托管式 Airflow(第 3 代) | 托管式 Airflow(第 2 代) | 托管式 Airflow(旧版第 1 代)
本教程介绍如何使用 Managed Service for Apache Airflow 创建 Apache Airflow DAG(有向无环图),该 DAG 可在 Managed Service for Apache Spark 集群上运行 Apache Hadoop WordCount 作业。
目标
- 访问您的 Managed Airflow 环境并使用 Airflow 界面。
- 创建和查看 Airflow 环境变量。
- 创建和运行包含以下任务的 DAG:
- 创建 Managed Service for Apache Spark 集群。
- 在该集群上运行 Apache Hadoop WordCount 作业。
- 将 WordCount 结果输出到 Cloud Storage 存储桶。
- 删除集群。
费用
在本文档中,您将使用 Google Cloud的以下收费组件:
- Managed Airflow
- Managed Service for Apache Spark
- Cloud Storage
如需根据您的预计使用情况来估算费用,请使用价格计算器。
准备工作
确保您的项目已启用以下 API:
控制台
启用 Managed Service for Apache Spark、Cloud Storage API。
启用 API 所需的角色
如需启用 API,您需要拥有
serviceusage.services.enable权限。如果您创建了项目,则可能已经通过 Owner 角色 (roles/owner) 获得了此权限。否则,您可以通过 Service Usage Admin 角色 (roles/serviceusage.serviceUsageAdmin) 获得此权限。了解如何授予角色。gcloud
启用 Managed Service for Apache Spark、Cloud Storage API:
启用 API 所需的角色
如需启用 API,您需要拥有
serviceusage.services.enable权限。如果您创建了项目,则可能已经通过 Owner 角色 (roles/owner) 获得了此权限。否则,您可以通过 Service Usage Admin 角色 (roles/serviceusage.serviceUsageAdmin) 获得此权限。了解如何授予角色。gcloud services enable dataproc.googleapis.com
storage-component.googleapis.com 在您的项目中,创建任意存储类别的 Cloud Storage 存储桶和区域,以用于存储 Hadoop WordCount 作业的结果。
请记住您所创建的存储桶的路径,例如
gs://example-bucket。您将为此路径定义一个 Airflow 变量,并在本教程稍后的示例 DAG 中使用该变量。使用默认参数创建 Managed Airflow 环境。等待环境创建操作完成。完成后,相应环境名称左侧会显示绿色对勾标记。
请记住您在其中创建了环境的区域,例如
us-central。您将为此区域定义 Airflow 变量,并在示例 DAG 中使用该变量,以在同一区域运行 Managed Service for Apache Spark 集群。
设置 Airflow 变量
设置 Airflow 变量,以便稍后在示例 DAG 中使用。例如,您可以在 Airflow 界面中设置 Airflow 变量。
| Airflow 变量 | 值 |
|---|---|
gcp_project
|
您在本教程中使用的项目的项目 ID,例如 example-project。 |
gcs_bucket
|
您为本教程创建的 Cloud Storage 存储桶的 URI,例如 gs://example-bucket |
gce_region
|
您创建环境的区域,例如 us-central1。
这是将创建 Managed Service for Apache Spark 集群的区域。 |
查看示例工作流
Airflow DAG 是您要安排和运行的有序任务的集合。DAG 在标准 Python 文件中定义。hadoop_tutorial.py 中显示的代码就是工作流代码。
运算符
为了编排示例工作流程中的三个任务,DAG 导入了以下三个 Airflow 运算符:
DataprocClusterCreateOperator:创建 Managed Service for Apache Spark 集群。DataProcHadoopOperator:用于提交 Hadoop Wordcount 作业并将结果写入 Cloud Storage 存储桶中。DataprocClusterDeleteOperator:用于删除集群,以避免产生持续的 Compute Engine 费用。
依赖项
您可以采用一种能够反映任务关系和依赖项的方式来组织要运行的任务。此 DAG 中的任务会按顺序执行。
时间安排
此 DAG 的名称为 composer_hadoop_tutorial,它会每天运行一次。由于传递到 default_dag_args 的 start_date 设置为 yesterday,因此 Managed Airflow 会安排工作流在 DAG 上传到环境的存储桶后立即开始运行。
将 DAG 上传到环境的存储桶
Managed Airflow 会将 DAG 存储在环境存储桶的 /dags 文件夹中。
如需上传 DAG,请执行以下操作:
在本地机器上,保存
hadoop_tutorial.py。在 Google Cloud 控制台中,前往环境页面。
在环境列表中,在您的环境对应的 DAGs 文件夹列中,点击 DAGs 链接。
点击上传文件。
在本地机器上选择
hadoop_tutorial.py,然后点击打开。
Managed Airflow 会将此 DAG 添加到 Airflow,并自动安排此 DAG。DAG 会在 3 到 5 分钟内发生更改。
探索 DAG 运行
查看任务状态
当您将 DAG 文件上传到 Cloud Storage 中的 dags/ 文件夹时,Managed Airflow 会解析该文件。成功完成后,工作流的名称会显示在 DAG 列表中,并且工作流会排入队列以立即运行。
如需查看任务状态,请转到 Airflow 网页界面并点击工具栏中的 DAGs。
如需打开 DAG 详情页面,请点击
composer_hadoop_tutorial。此页面包含工作流任务和依赖项的图形表示。
如需查看每个任务的状态,请点击 Graph View,然后将鼠标悬停在每个任务对应的图形上。
将工作流重新加入队列
如需从 Graph View 重新运行工作流,请执行以下操作:
- 在 Airflow 界面的“Graph View”中,点击
create_dataproc_cluster图形。 - 如需重置三个任务,请点击 Clear,然后点击 OK 进行确认。
- 在“Graph View”中,再次点击
create_dataproc_cluster。 - 如需将工作流重新加入队列,请点击 Run。
查看任务结果
您还可以前往以下 Google Cloud 控制台页面来检查 composer_hadoop_tutorial 工作流的状态和结果:
Managed Service for Apache Spark 集群:用于监控集群创建和删除。请注意,由工作流创建的集群是临时性的,也就是说,此类集群仅在工作流的持续期间内存在,并且将在最后一个工作流任务的执行过程中删除。
Managed Service for Apache Spark 作业:可以查看或监控 Apache Hadoop Wordcount 作业。点击“任务 ID”即可查看作业日志输出。
Cloud Storage 浏览器:可以在您为本教程创建的 Cloud Storage 存储桶所含的
wordcount文件夹中查看 WordCount 的结果。
清理
删除本教程中使用的资源:
删除 Managed Airflow 环境,包括手动删除环境的存储桶。
删除 Cloud Storage 存储桶,该存储分区用于存储 Hadoop WordCount 作业的结果。