ARTICLE DETAIL

资讯详情

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

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

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多节点模式运行包含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-Reqs1. 标准化 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 写入.envmkdir -p ./dags ./logs ./plugins echo -e AIRFLOW_UID$(id -u) .envWindows 用户同样建议执行上述命令若使用 MINGW/GitBash命令完全一致。如果不想在启动日志中看到AIRFLOW_UID is not set告警也可以直接手动创建.env并写入AIRFLOW_UID50000这个.env会被 docker-compose 自动读取对应到 docker-compose.yaml 中的user: ${AIRFLOW_UID:-50000}:0配置——注意官方模板把容器内用户组固定为0root 组以保证各容器间对挂载卷的读写一致。第 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而我们还要安装gcloudGoogle 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 模板版本对应首次构建约需 515 分钟视网络而定。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与 Hookpyarrow提供 CSV → Parquet 格式转换能力对应 DAG 中的format_to_parquet函数。SHELL [/bin/bash, -o, pipefail, -e, -u, -x, -c] ARG CLOUD_SDK_VERSION322.0.0 ENV GCLOUD_HOME/home/google-cloud-sdk ENV PATH${GCLOUD_HOME}/bin/:${PATH} RUN DOWNLOAD_URLhttps://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-components1 \ ${GCLOUD_HOME}/install.sh \ --bash-completionfalse \ --path-updatefalse \ --usage-reportingfalse \ --quiet \ rm -rf ${TMP_DIR} \ gcloud --version这一段是镜像定制的核心下载并安装google-cloud-sdk版本 322.0.0Linux 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标准 ADCApplication Default Credentials路径让google.cloudPython 客户端自动读取凭据必须指向容器内路径/.google/credentials/google_credentials.jsonAIRFLOW_CONN_GOOGLE_CLOUD_DEFAULT定义 Airflow 默认 GCP 连接Connection供 GCP provider 的 operator 使用以google-cloud-platform://为 URI通过extra__google_cloud_platform__key_path指定凭据文件GCP_PROJECT_IDGCP 项目 IDDAG 中通过os.environ.get(GCP_PROJECT_ID)读取如pivotal-surfer-336713GCP_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保证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 fhttps://s3.amazonaws.com/nyc-tlc/tripdata/{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_taskBashOperatorcurl -sSL dataset_url下载 CSV 到 Airflow 家目录format_to_parquet_taskPythonOperator调用format_to_parquet用pyarrow.csv读入 CSV、pyarrow.parquet写出 Parquet——对应 requirements.txt 中pyarrow的角色local_to_gcs_taskPythonOperator调用upload_to_gcs上传raw/文件名.parquet到 GCS 桶。函数内对storage.blob的 multipart 阈值与 chunk size 做了 5MB 的 workaround 调整规避慢速上传大文件时的超时问题bigquery_external_table_taskBigQueryCreateExternalTableOperator来自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 组件的容器 IDdocker 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),仅供参考
返回列表