本文属于机器翻译版本。若本译文内容与英语原文存在差异,则一律以英文原文为准。
在 Amazon MWAA 中使用 dbt
借助亚马逊 MWAA,您可以使用 dbt(数据构建工具)和 PostgreSQL 来构建和运行数据转换工作流程。在以下步骤中,使用启动脚本添加所需的依赖关系,并将示例 dbt 项目上传到环境的 Amazon S3 存储桶。然后,使用示例 DAG 来验证亚马逊 MWAA 是否已安装依赖项。最后,使用运行 BashOperator dbt 项目。
版本
你可以使用本页上的代码示例,在 Python 网站上使用 Python 3.12 中的 Apache Airflow v2,
先决条件
在完成以下步骤之前,您需要满足以下条件:
-
使用 Apache Airflow v2.11.2 的亚马逊 MWAA 环境。此示例是使用 v2.11.2 编写和测试的。您可能需要修改示例以与其他 Apache Airflow 版本一起使用。
-
dbt 项目示例。要开始在亚马逊 MWAA 中使用 dbt,您可以创建一个分叉并从 dbt-labs 存储库中克隆
dbt 入门项目。 GitHub
依赖项
要将 Amazon MWAA 与 dbt 配合使用,请将以下启动脚本添加到环境中。要了解更多信息,请参阅在 Amazon MWAA 中使用启动脚本。
#!/bin/bash if [[ "${MWAA_AIRFLOW_COMPONENT}" != "worker" ]] then exit 0 fi echo "------------------------------" echo "Installing virtual Python env" echo "------------------------------" pip3 install --upgrade pip echo "Current Python version:" python3 --version echo "..." sudo pip3 install --user virtualenv sudo mkdir -p /usr/local/airflow/python3-virtualenv cd /usr/local/airflow/python3-virtualenv sudo python3 -m venv dbt-env sudo chmod -R 777 * echo "------------------------------" echo "Activating venv in $DBT_ENV_PATH" echo "------------------------------" source dbt-env/bin/activate pip3 list echo "------------------------------" echo "Installing libraries..." echo "------------------------------" # do not use sudo, as it will install outside the venv pip3 install dbt-core==1.9.4 dbt-redshift==1.9.1 dbt-postgres==1.9.0 echo "------------------------------" echo "Venv libraries..." echo "------------------------------" pip3 list dbt --version echo "------------------------------" echo "Deactivating venv..." echo "------------------------------" deactivate
设置 DBT_ENV_PATH 变量
您可以在启动脚本$DBT_ENV_PATH中进行设置,也可以在亚马逊 MWAA 环境中将其设置为 Airflow 配置。
在以下部分中,将您的 dbt 项目目录上传到 Amazon S3 并运行 DAG 来验证亚马逊 MWAA 是否已成功安装所需的债务依赖项。
将 dbt 项目上传到 Amazon S3
为了能够在 Amazon MWAA 环境中使用 dbt 项目,您可以将整个项目目录上传到环境的 dags 文件夹中。当环境更新时,Amazon MWAA 会将 dbt 目录下载到本地 usr/local/airflow/dags/ 文件夹。
要将 dbt 项目上传到 Amazon S3,请执行以下操作
-
导航到您克隆 dbt 入门项目的目录。
-
运行以下 Amazon S3 AWS CLI 命令,使用
--recursive参数以递归方式将项目内容复制到您的环境dags文件夹。该命令会创建一个名为dbt的子目录,您可以将其用于所有 dbt 项目。如果子目录已经存在,则项目文件将被复制到现有目录中,并且不会创建新目录。该命令还会为该特定入门项目在dbt目录中创建一个子目录。aws s3 cpdbt-starter-projects3://amzn-s3-demo-bucket/dags/dbt/dbt-starter-project--recursive您可以为项目子目录使用不同的名称,以便在
dbt父目录中组织多个 dbt 项目。
使用 DAG 验证 dbt 依赖项的安装
以下 DAG 使用BashOperator和 bash 命令来验证 Amazon MWAA 是否已成功安装启动脚本中指定的 dbt 依赖项。
from airflow import DAG from airflow.operators.bash_operator import BashOperator from airflow.utils.dates import days_ago with DAG(dag_id="dbt-installation-test", schedule_interval=None, catchup=False, start_date=days_ago(1)) as dag: cli_command = BashOperator( task_id="bash_command", bash_command="/usr/local/airflow/python3-virtualenv/dbt-env/bin/dbt --version" )
执行以下操作以访问任务日志并验证是否已安装 dbt 及其依赖项。
-
导航到 Amazon MWAA 控制台,然后从可用环境列表中选择打开 Airflow UI。
-
在 Apache Airflow UI 上,从列表中找到
dbt-installation-testDAG,然后在Last Run列中选择日期以打开上一个成功的任务。 -
使用图表视图,选择
bash_command任务以打开任务实例的详细信息。 -
选择 “记录” 打开任务日志,然后验证日志是否成功列出了启动脚本中指定的 dbt 版本。
创建并上传 dbt profiles.yml
要连接到目标数据库,dbt 需要一个profiles.yml文件。下一节中的 DAG 通过--profiles-dir /tmp/dbt,所以 dbt profiles.yml 直接在/tmp/dbt目录中查找。这是您上传到亚马逊 S3 的dbt文件夹。
中的配置文件名称profiles.yml必须与启动项目中定义的profile:值相匹配dbt_project.yml。dbt 入门项目使用default.
要创建和上传 profiles.yml
-
在克隆启动项目的目录中,
profiles.yml使用您的数据库连接详细信息创建一个名为的文件。default: target: dev outputs: dev: type: postgres host:your-db-endpoint.region.rds.amazonaws.com port: 5432 user:your_db_userpassword:your_db_passworddbname:your_databaseschema:your_schemathreads: 4 -
将文件上传到环境的 DAGs
dbt文件夹中的子目录。aws s3 cp profiles.yml s3://amzn-s3-demo-bucket/dags/dbt/profiles.yml
保护数据库凭据
为避免存储纯文本凭证,请使用 dbt env_var() 函数引用机密。通过亚马逊 MWAA 环境变量或 AWS Secrets Manager——例如,提供值。password: "{{ env_var('DBT_PASSWORD') }}"还要确保您的 Amazon MWAA VPC 安全组允许您的工作人员通过配置的端口访问您的数据库。
使用 DAG 来运行 dbt 项目
以下 DAG 使用 BashOperator 将您从本地 usr/local/airflow/dags/ 目录上传到 Amazon S3 的 dbt 项目复制到可写入的 /tmp 目录,然后运行 dbt 项目。bash 命令假设一个名为 dbt-starter-project 的 入门 dbt 项目。根据您项目目录的名称修改目录名称。
from airflow import DAG from airflow.operators.bash_operator import BashOperator from airflow.utils.dates import days_ago import os DAG_ID = os.path.basename(__file__).replace(".py", "") # assumes all files are in a subfolder of DAGs called dbt with DAG(dag_id=DAG_ID, schedule_interval=None, catchup=False, start_date=days_ago(1)) as dag: cli_command = BashOperator( task_id="bash_command", bash_command="source /usr/local/airflow/python3-virtualenv/dbt-env/bin/activate;\ cp -R /usr/local/airflow/dags/dbt /tmp;\ echo 'listing project files:';\ ls -R /tmp;\ cd /tmp/dbt/dbt-starter-project;\ /usr/local/airflow/python3-virtualenv/dbt-env/bin/dbt run --project-dir /tmp/dbt/dbt-starter-project --profiles-dir /tmp/dbt;\ cat /tmp/dbt/dbt-starter-project/logs/dbt.log;\ rm -rf /tmp/dbt/dbt-starter-project" )