Data Engineering Zoomcamp:在 Docker 中以官方方式搭建 Apache Airflow,编排 NYC Taxi 数据摄入 GCS/BigQuery 管线
【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 👇🏼项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp
本文是 Data Engineering Zoomcamp 第二周「工作流编排」模块的实战指南,聚焦于使用 Airflow 官方提供的 docker-compose 模板,在本地 Docker 环境中搭建一套多组件(CeleryExecutor + Redis + PostgreSQL)的 Airflow 集群,并通过定制 Dockerfile 与挂载 GCP 服务账号凭据,打通「下载 NYC Taxi 数据 → 转 Parquet → 上传 GCS → 建 BigQuery 外部表」的完整数据摄入管线。读完本文,你将掌握官方方式 Airflow 环境的标准落地流程、关键配置项的含义,以及常见凭据挂载故障的排查手段。
一、方案概览:官方模板 vs 轻量定制
本模块提供两条 Airflow 本地部署路线(见 airflow/README.md):
- 官方版(本文主体):直接拉取 Airflow 官方 docker-compose.yaml,以
CeleryExecutor多节点模式运行,包含postgres、redis、airflow-webserver、airflow-scheduler、airflow-worker、airflow-triggerer、airflow-init等 7 个服务; - No-Frills 轻量版:基于 2_setup_nofrills.md 的说明,精简为
LocalExecutor单节点模式,仅保留postgres、scheduler、webserver三个服务,内存占用更低。
官方模板服务众多乍看令人望而生畏,但它只是一个 quick-start 模板,随着你对各组件职责的理解加深,完全可以逐步移除用不到的服务——本仓库中的 docker-compose-nofrills.yml 就是这个模板的「去繁就简」版本范例。
二、前置准备(Pre-Reqs)
1. 标准化 GCP 服务账号凭据
为保证本 workshop 中所有配置的一致性,需要把 GCP 服务账号的凭据文件统一重命名为google_credentials.json,并存放于$HOME目录下的固定位置:
cd ~ && mkdir -p ~/.google/credentials/ mv <path/to/your/service-account-authkeys>.json ~/.google/credentials/google_credentials.json注意:凭据路径、文件名必须与后文 docker-compose 的
volumes挂载、GOOGLE_APPLICATION_CREDENTIALS环境变量严格对齐,任何一处不一致都会导致「凭据文件找不到」的经典报错(详见第七节排查)。
2. Docker 环境要求
- docker-compose 版本:建议升级到 v2.x 以上;
- Docker Engine 内存:至少分配 5GB(官方版),理想值 8GB。官方模板的
airflow-init容器在启动时会校验可用资源(内存 ≥ 4GB、CPU ≥ 2 核、磁盘 ≥ 10GB),内存不足会打印显式告警; - 若分配内存不足,最典型的现象是airflow-webserver 容器持续重启(restart 循环),因为 webserver 启动时加载全部 DAG 与元数据会吃满内存。
3. Python 版本
宿主机 Python 版本要求 3.7+(用于在本地调试 DAG 脚本;容器内 Airflow 自身的 Python 环境由镜像自带,不受宿主机影响)。
三、Airflow 官方方式安装步骤
第 1 步:创建项目子目录
在项目根目录下新建airflow子目录(即本仓库cohorts/2022/week_2_data_ingestion/airflow/目录的对应物),后续所有文件都在此目录内操作。
第 2 步:设置 Airflow 用户(Linux 关键步骤)
Airflow 官方镜像内的默认用户 UID 是 50000。在 Linux 上,如果不显式声明AIRFLOW_UID,容器内以 root 身份创建的文件(dags、logs、plugins三个目录)会归 root 所有,宿主机普通用户将无法读写,导致日志无法落盘、DAG 无法编辑。
标准做法是创建好三个挂载目录,并把当前用户 UID 写入.env:
mkdir -p ./dags ./logs ./plugins echo -e "AIRFLOW_UID=$(id -u)" > .envWindows 用户同样建议执行上述命令;若使用 MINGW/GitBash,命令完全一致。如果不想在启动日志中看到AIRFLOW_UID is not set告警,也可以直接手动创建.env并写入:
AIRFLOW_UID=50000这个.env会被 docker-compose 自动读取,对应到 docker-compose.yaml 中的user: "${AIRFLOW_UID:-50000}:0"配置——注意官方模板把容器内用户组固定为0(root 组),以保证各容器间对挂载卷的读写一致。
第 3 步:拉取官方 docker-compose 模板
从最新版 Apache Airflow 官方文档下载 docker-compose 文件:
curl -LfO 'https://airflow.apache.org/docs/apache-airflow/stable/docker-compose.yaml'仓库内已保存了一份基于 Airflow 2.2.3 的成品:docker-compose.yaml,可直接对照或复制使用。
下载下来的模板包含 7 个服务,初次接触会觉得「过度设计」。请理解:这是官方为了覆盖CeleryExecutor完整集群形态给出的通用 quick-start,其中redis充当 Celery 的 broker(消息队列)、postgres既作元数据库又作 Celery result backend、flower提供 Celery 监控 UI、airflow-triggerer服务 2.2+ 的延迟触发任务。实际使用时可按需裁剪。
第 4 步:定制 Dockerfile 构建扩展镜像
为什么需要扩展镜像?因为官方基础镜像只包含 Airflow 核心与默认 provider,而我们还要:
- 安装
gcloud(Google Cloud SDK),用于连接 GCS 数据湖(Data Lake); - 通过
requirements.txt用pip install安装额外的 Python 依赖。
仓库中的成品 Dockerfile 完整呈现了这套定制逻辑,逐段解读如下:
# First-time build can take upto 10 mins. FROM apache/airflow:2.2.3 ENV AIRFLOW_HOME=/opt/airflow以apache/airflow:2.2.3为基础镜像(与下载的 docker-compose 模板版本对应),首次构建约需 5~15 分钟(视网络而定)。
USER root RUN apt-get update -qq && apt-get install vim -qqq COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt切换为 root 用户安装系统级工具vim(便于进容器改文件),并把 requirements.txt 复制进镜像执行 pip 安装。本仓库的 requirements 只有两行,均为 GCS 摄入链路的关键依赖:
apache-airflow-providers-google pyarrowapache-airflow-providers-google:提供 GCP 相关 operator(如BigQueryCreateExternalTableOperator)与 Hook;pyarrow:提供 CSV → Parquet 格式转换能力(对应 DAG 中的format_to_parquet函数)。
SHELL ["/bin/bash", "-o", "pipefail", "-e", "-u", "-x", "-c"] ARG CLOUD_SDK_VERSION=322.0.0 ENV GCLOUD_HOME=/home/google-cloud-sdk ENV PATH="${GCLOUD_HOME}/bin/:${PATH}" RUN DOWNLOAD_URL="https://dl.google.com/dl/cloudsdk/channels/rapid/downloads/google-cloud-sdk-${CLOUD_SDK_VERSION}-linux-x86_64.tar.gz" \ && TMP_DIR="$(mktemp -d)" \ && curl -fL "${DOWNLOAD_URL}" --output "${TMP_DIR}/google-cloud-sdk.tar.gz" \ && mkdir -p "${GCLOUD_HOME}" \ && tar xzf "${TMP_DIR}/google-cloud-sdk.tar.gz" -C "${GCLOUD_HOME}" --strip-components=1 \ && "${GCLOUD_HOME}/install.sh" \ --bash-completion=false \ --path-update=false \ --usage-reporting=false \ --quiet \ && rm -rf "${TMP_DIR}" \ && gcloud --version这一段是镜像定制的核心:下载并安装google-cloud-sdk(版本 322.0.0,Linux x86_64),通过install.sh静默安装,关闭 bash 补全、路径更新与用量上报,最后用gcloud --version验证安装成功。安装后的gcloud二进制位于$GCLOUD_HOME/bin/,已加入PATH,容器内即可直接执行 gsutil/gcloud 命令。
WORKDIR $AIRFLOW_HOME COPY scripts scripts RUN chmod +x scripts USER $AIRFLOW_UID回到$AIRFLOW_HOME(/opt/airflow),复制脚本目录并授予执行权限,最后切回$AIRFLOW_UID用户运行(与第 2 步的 UID 设置呼应,避免容器内以 root 运行产生文件属主问题)。
第 5 步:改造 docker-compose.yaml
回到下载的docker-compose.yaml,针对x-airflow-common锚点(被所有 airflow 组件服务共享的基础配置块)做三处修改,仓库成品 docker-compose.yaml 即改造结果:
① 用 build 替换 image
注释或删除image标签,改为从本地 Dockerfile 构建(这样第 4 步的扩展依赖才生效):
x-airflow-common: &airflow-common build: context: . dockerfile: ./Dockerfile② 挂载 GCP 凭据为只读
在volumes中把宿主机~/.google/credentials/目录挂载到容器内固定路径/.google/credentials,并加上:ro(只读)防止容器内误改凭据:
volumes: - ./dags:/opt/airflow/dags - ./logs:/opt/airflow/logs - ./plugins:/opt/airflow/plugins - ~/.google/credentials/:/.google/credentials:ro③ 注入 GCP 相关环境变量
在environment中设置四个关键变量(仓库成品中已给出示例值,TODO注释提示按你自己的 GCP 配置替换):
environment: GOOGLE_APPLICATION_CREDENTIALS: /.google/credentials/google_credentials.json AIRFLOW_CONN_GOOGLE_CLOUD_DEFAULT: 'google-cloud-platform://?extra__google_cloud_platform__key_path=/.google/credentials/google_credentials.json' # TODO: Please change GCP_PROJECT_ID & GCP_GCS_BUCKET, as per your config GCP_PROJECT_ID: 'pivotal-surfer-336713' GCP_GCS_BUCKET: 'dtc_data_lake_pivotal-surfer-336713'这四个变量的作用分别是:
| 环境变量 | 作用 | 取值要点 |
|---|---|---|
GOOGLE_APPLICATION_CREDENTIALS | 标准 ADC(Application Default Credentials)路径,让google.cloudPython 客户端自动读取凭据 | 必须指向容器内路径/.google/credentials/google_credentials.json |
AIRFLOW_CONN_GOOGLE_CLOUD_DEFAULT | 定义 Airflow 默认 GCP 连接(Connection),供 GCP provider 的 operator 使用 | 以google-cloud-platform://为 URI,通过extra__google_cloud_platform__key_path指定凭据文件 |
GCP_PROJECT_ID | GCP 项目 ID,DAG 中通过os.environ.get("GCP_PROJECT_ID")读取 | 如pivotal-surfer-336713 |
GCP_GCS_BUCKET | GCS 数据湖桶名,DAG 上传与外部表 URI 的来源 | 命名规范常为dtc_data_lake_<project-id> |
④(可选)关闭示例 DAG
将AIRFLOW__CORE__LOAD_EXAMPLES改为false,避免每次启动加载 Airflow 自带的 40+ 个示例 DAG,减少 UI 噪音与启动时间(仓库成品已默认关闭)。
四、启动集群与验证
按 airflow/README.md 的执行清单操作:
# 1. 构建镜像(首次构建约 15 分钟,之后仅当 Dockerfile/requirements 变化时才需重建) docker-compose build # 2. 初始化:升级元数据库、创建管理员账号(airflow-init 一次性服务) docker-compose up airflow-init # 3. 启动全部服务 docker-compose up # 4. 另开终端,确认 7 个容器全部 Up docker-compose ps # 5. 浏览器访问 http://localhost:8080 ,默认账号密码 airflow/airflow各服务的启动顺序由depends_on与健康检查(healthcheck)保证:postgres、redis先就绪,airflow-init完成数据库升级与 admin 创建,随后 webserver/scheduler/worker 才正式接管任务。若中途看到AIRFLOW_UID not set或资源不足告警,属于提示级别,不影响启动;但资源严重不足时 webserver 会陷入「健康检查失败 → restart」的循环,此时应回到第二节加大 Docker 内存配额。
五、用官方栈跑通第一条数据摄入 DAG
环境就绪后,把 DAG 脚本放入挂载的./dags目录即可被调度器自动发现。仓库中 dags/data_ingestion_gcs_dag.py 是一条完整的 NYC Taxi 数据摄入流水线,其四个任务与上图 Graph 视图一一对应:
PROJECT_ID = os.environ.get("GCP_PROJECT_ID") BUCKET = os.environ.get("GCP_GCS_BUCKET") dataset_file = "yellow_tripdata_2021-01.csv" dataset_url = f"https://s3.amazonaws.com/nyc-tlc/trip+data/{dataset_file}" path_to_local_home = os.environ.get("AIRFLOW_HOME", "/opt/airflow/") parquet_file = dataset_file.replace('.csv', '.parquet') BIGQUERY_DATASET = os.environ.get("BIGQUERY_DATASET", 'trips_data_all')GCP_PROJECT_ID、GCP_GCS_BUCKET正是上一节注入的环境变量,DAG 通过os.environ.get(...)读取——这正是官方 docker-compose 配置「环境变量驱动」设计的意义所在;- 数据集源为 S3 上的 NYC TLC 公开数据(
yellow_tripdata_2021-01.csv)。
四个任务依次为:
- download_dataset_task(
BashOperator):curl -sSL <dataset_url>下载 CSV 到 Airflow 家目录; - format_to_parquet_task(
PythonOperator):调用format_to_parquet,用pyarrow.csv读入 CSV、pyarrow.parquet写出 Parquet——对应 requirements.txt 中pyarrow的角色; - local_to_gcs_task(
PythonOperator):调用upload_to_gcs上传raw/<文件名>.parquet到 GCS 桶。函数内对storage.blob的 multipart 阈值与 chunk size 做了 5MB 的 workaround 调整,规避慢速上传大文件时的超时问题; - bigquery_external_table_task(
BigQueryCreateExternalTableOperator,来自apache-airflow-providers-google):以gs://<bucket>/raw/<文件名>.parquet为sourceUris、PARQUET为sourceFormat创建 BigQuery 外部表,指向trips_data_all.external_table。
任务依赖链:
download_dataset_task >> format_to_parquet_task >> local_to_gcs_task >> bigquery_external_table_task在 Airflow Web UI 中触发该 DAG 后,即可在 Graph 视图看到上图所示的状态流转。至此,官方栈的价值完整闭环:Docker 提供隔离一致的运行时,Airflow 提供可调度、可重试、可视化监控的编排能力,GCP 侧完成数据湖(GCS)+ 数仓(BigQuery)的落地。
六、两种备选方案速览
官方模板之外,本仓库还提供两种轻量替代:
- No-Frills 轻量版:参见 2_setup_nofrills.md 与 docker-compose-nofrills.yml。相比官方版的主要差异:移除
redis/worker/triggerer/flower/airflow-init五个服务,执行器从CeleryExecutor(多节点)改为LocalExecutor(单节点);改用.env集中管理变量;新增 scripts/entrypoint.sh 作为 webserver 入口,负责airflow db upgrade与创建admin/admin用户。切换前需执行docker-compose down --volumes --rmi all清理旧环境。 - 2.3.4 轻量本地版:仓库另存有 docker-compose_2.3.4.yaml,可重命名为
docker-compose.yaml直接使用,同样需替换GCP_PROJECT_ID与GCP_GCS_BUCKET。
七、常见问题排查
File /.google/credentials/google_credentials.json was not found
这是本 workshop 最高频的报错,按以下顺序排查:
第一步:确认宿主机文件存在且命名正确
凭据必须位于$HOME/.google/credentials/,且文件名必须是google_credentials.json(对应第二节第 1 步的标准化动作)。
第二步:确认容器内确实能看到该文件
查看正在运行的容器,找到 airflow worker(或任意 airflow 组件)的容器 ID:
docker ps进入容器:
docker exec -it <container-ID> bash在容器内检查凭据目录:
ls -lh /.google/credentials/第三步:若目录为空,改用绝对路径挂载
如果容器内该目录为空,说明 docker-compose 没能把宿主机目录映射进去(多为~展开或路径解析问题)。此时把volumes中该行改为宿主机绝对路径,例如 Windows 下的写法:
volumes: - ./dags:/opt/airflow/dags - ./logs:/opt/airflow/logs - ./plugins:/opt/airflow/plugins # here: ---------------------------- - c:/Users/alexe/.google/credentials/:/.google/credentials:ro # -----------------------------------改完重新docker-compose up -d并再次docker exec验证。
其他注意事项
- Windows/WSL 用户在 No-Frills 方案下若遇到
ModuleNotFoundError或性能问题,请核对 WSL2 环境下 Docker 的资源与挂载配置; - 凭据路径若有自定义(如放在
$HOME/.gc),需同步修改三处:GOOGLE_APPLICATION_CREDENTIALS、AIRFLOW_CONN_GOOGLE_CLOUD_DEFAULT以及volumes挂载路径,三者必须指向同一个文件。
八、清理与后续
运行完毕或需要切换方案时,按需选择清理力度:
# 停止并删除容器(保留数据卷与镜像) docker-compose down # 停止并删除容器、清空数据卷、删除镜像(完全重置) docker-compose down --volumes --rmi all # 或仅清理孤儿容器/数据卷 docker-compose down --volumes --remove-orphans若长期不清理 Docker 缓存,也可执行docker system prune释放空间,同时清空 airflow 的logs目录。
延伸阅读
- Airflow 概念与架构详解:docs/1_concepts.md
- 轻量级 No-Frills 安装说明:2_setup_nofrills.md
- 完整执行清单与两个版本的对照:airflow/README.md
- 官方文档中与本方案相关的主题包括:Docker 快速开始、Docker 镜像构建(docker-stack/build)与镜像定制配方(docker-stack/recipes),可在 Apache Airflow 官网对应页面查阅,本文的 Dockerfile 即参考 recipes 思路定制而成。
【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 👇🏼项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考