news 2026/9/11 14:00:13

Data Engineering Zoomcamp:在 Docker 中以官方方式搭建 Apache Airflow,编排 NYC Taxi 数据摄入 GCS/BigQuery 管线

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Data Engineering Zoomcamp:在 Docker 中以官方方式搭建 Apache Airflow,编排 NYC Taxi 数据摄入 GCS/BigQuery 管线

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多节点模式运行,包含postgresredisairflow-webserverairflow-schedulerairflow-workerairflow-triggererairflow-init等 7 个服务;
  • No-Frills 轻量版:基于 2_setup_nofrills.md 的说明,精简为LocalExecutor单节点模式,仅保留postgresschedulerwebserver三个服务,内存占用更低。

官方模板服务众多乍看令人望而生畏,但它只是一个 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 身份创建的文件(dagslogsplugins三个目录)会归 root 所有,宿主机普通用户将无法读写,导致日志无法落盘、DAG 无法编辑。

标准做法是创建好三个挂载目录,并把当前用户 UID 写入.env

mkdir -p ./dags ./logs ./plugins echo -e "AIRFLOW_UID=$(id -u)" > .env

Windows 用户同样建议执行上述命令;若使用 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,而我们还要:

  1. 安装gcloud(Google Cloud SDK),用于连接 GCS 数据湖(Data Lake);
  2. 通过requirements.txtpip 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 pyarrow
  • apache-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_IDGCP 项目 ID,DAG 中通过os.environ.get("GCP_PROJECT_ID")读取pivotal-surfer-336713
GCP_GCS_BUCKETGCS 数据湖桶名,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)保证:postgresredis先就绪,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_IDGCP_GCS_BUCKET正是上一节注入的环境变量,DAG 通过os.environ.get(...)读取——这正是官方 docker-compose 配置「环境变量驱动」设计的意义所在;
  • 数据集源为 S3 上的 NYC TLC 公开数据(yellow_tripdata_2021-01.csv)。

四个任务依次为:

  1. download_dataset_taskBashOperator):curl -sSL <dataset_url>下载 CSV 到 Airflow 家目录;
  2. format_to_parquet_taskPythonOperator):调用format_to_parquet,用pyarrow.csv读入 CSV、pyarrow.parquet写出 Parquet——对应 requirements.txt 中pyarrow的角色;
  3. local_to_gcs_taskPythonOperator):调用upload_to_gcs上传raw/<文件名>.parquet到 GCS 桶。函数内对storage.blob的 multipart 阈值与 chunk size 做了 5MB 的 workaround 调整,规避慢速上传大文件时的超时问题;
  4. bigquery_external_table_taskBigQueryCreateExternalTableOperator,来自apache-airflow-providers-google):以gs://<bucket>/raw/<文件名>.parquetsourceUrisPARQUETsourceFormat创建 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_IDGCP_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_CREDENTIALSAIRFLOW_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/11 13:59:31

老旧数传电台升级Orbit LN:工业无线专网技改实测与决策指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/11 13:59:15

Windows上运行Linux的5种高效方案对比

1. 为什么要在Windows上运行Linux&#xff1f;作为一名在Windows和Linux双系统间反复横跳多年的开发者&#xff0c;我深刻理解这种切换带来的痛苦。每次需要编译Linux程序就得重启进入Ubuntu&#xff0c;调试完再切回Windows处理日常工作&#xff0c;这种割裂感严重影响效率。直…

作者头像 李华
网站建设 2026/9/11 13:54:59

RT-Thread嵌入式AI工业质检实战:低代码部署YOLOv5s

1. 这不是“跑个Demo”的AI&#xff0c;而是嵌入式工程师能亲手焊进产线的工业质检系统你有没有见过这样的场景&#xff1a;工厂车间里&#xff0c;一台老式PLC控制的传送带正把刚压铸出来的金属支架送进检测工位。旁边站着两位老师傅&#xff0c;一人盯着流水线&#xff0c;一…

作者头像 李华
网站建设 2026/9/11 13:53:47

群体PCA分析:核心价值、应用场景与实战技巧

1. 群体PCA分析的核心价值与应用场景 主成分分析&#xff08;PCA&#xff09;作为降维利器&#xff0c;在生物信息学、金融风控、消费行为研究等领域已成为标准分析流程。当样本量达到数百甚至上千时&#xff0c;传统的二维散点图已难以清晰展示群体结构特征。这时群体PCA分析就…

作者头像 李华