Data Engineering Zoomcamp 模块三:基于 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 第三模块「Data Warehouse and BigQuery」展开,系统讲解数据仓库的核心概念(OLTP/OLAP、数据仓库架构、Data Mart)、BigQuery 的无服务器特性与定价模型,并以纽约出租车数据为例,通过真实可运行的 SQL 演示外部表、分区(Partitioning)与聚簇(Clustering)的建表与优化效果,最后覆盖 BigQuery ML 的建模与模型部署全流程。读完本文,你将掌握用 BigQuery 从「建外部表 → 分区聚簇优化 → 成本与性能最佳实践 → SQL 内建机器学习模型 → Docker 部署模型」的完整数据仓库实战链路。
模块总览:学什么、怎么学
03-data-warehouse/README.md是模块三的入口文档,它把整个模块组织为「视频课程 + 可复现 SQL + 作业 + 社区笔记」四部分:
- 课程视频:共 6 讲,覆盖 Data Warehouse 与 BigQuery、Partitioning vs Clustering、Best practices、Internals of BigQuery,以及两个进阶主题——BigQuery 机器学习与模型部署;
- 可执行 SQL 文件:
big_query.sql(基础查询、外部表、分区聚簇)、big_query_ml.sql(机器学习建模)、big_query_hw.sql(作业配套 SQL); - 部署手册:
extract_model.md(从 BigQuery 导出模型并用 Docker 提供服务); - 作业:指向 2026 年模块三作业,同时配套数据加载脚本;
- 社区笔记:README 的 Community notes 区块收集了历届学员的公开笔记(详见
03-data-warehouse/README.md的 details 折叠区)。
此外,cohorts/2027/03-data-warehouse/下还提供了与视频逐讲对应的文字讲义(01-data-warehouse-and-bigquery.md至06-deploying-a-machine-learning-model.md),本指南将结合这些讲义与仓库源码,对 README 中的主题做纵深展开。
数据仓库基础:OLTP 与 OLAP 的区别
在动手使用 BigQuery 之前,需要先理解数据仓库在整个数据体系中的定位。数据库按用途分为两大流派:
- OLTP(Online Transaction Processing,联机事务处理):服务于后端业务系统。多个 SQL 语句打包在一个事务中执行,任一语句失败则整体回滚。典型例子是电商下单:订单、支付、库存更新要么全部成功、要么全部失败。
- OLAP(Online Analytical Processing,联机分析处理):服务于分析场景,核心目标是「把大量数据装进去,并从中发现隐藏洞察」,主要使用者是数据分析师与数据科学家。
两种数据库几乎在所有维度上都存在差异,二者的对比可用下表概括:
| 维度 | OLTP | OLAP |
|---|---|---|
| 目的 | 实时控制与支撑关键业务运营 | 规划、解决问题、辅助决策、发现隐藏洞察 |
| 数据更新 | 用户触发的短小快速更新 | 通过定时、长时批处理任务周期性刷新 |
| 数据库设计 | 为效率而规范化(Normalized) | 为分析而反规范化(Denormalized) |
| 空间需求 | 归档历史后通常较小 | 因聚合海量数据集通常很大 |
数据仓库架构:从数据源到 Data Mart
数据仓库(Data Warehouse)本质上就是一套面向报表与数据分析的 OLAP 方案。一套典型的仓库由原始数据(raw data)、元数据(metadata)和汇总数据(summary data)组成,其数据流如下:多个数据源(运营系统、平面文件、OLTP 数据库等)先汇入暂存区(staging area),再由暂存区写入数据仓库本体,最终按主题切分为更小、面向特定业务的Data Mart(如采购、销售、库存),供不同用户直接访问。
仓库的价值在于它的灵活性:对分析师而言,通过 Data Mart 访问数据是理想形态;对数据科学家而言,也可以直接读取仓库中的原始数据。
BigQuery:无服务器的数据仓库
BigQuery 是模块使用的数据仓库实现,其最大优势是无服务器(serverless):无需管理任何服务器,也无需安装数据库软件。企业自建数据仓库时,大量时间都耗在基础设施的创建与维护上,而 BigQuery 将软件与基础设施一并托管,开箱即用地提供可扩展性与高可用性——从几 GB 起步,可平滑扩展到 PB 级。
BigQuery 的差异化能力包括:
- 通过 SQL 接口直接做机器学习(即 BigQuery ML,见下文);
- 支持**地理空间数据(geospatial)**处理;
- 支持**商业智能(BI)**类查询。
另一个关键设计是存储与计算分离:传统架构中单台大服务器同时承载存储与计算,数据量增长时整机必须随之扩容;BigQuery 则将计算引擎与存储解耦,数据放在用户自选的存储(如 Cloud Storage)上,按需分析。这一设计显著降低了成本。
定价模型
BigQuery 提供两种定价方式:
- 按需计费(On-demand):按扫描/处理的数据量计费,每处理 1 TB 数据约 $5;
- 固定费率(Flat rate):按预购的 slot(BigQuery 的处理能力单位)计费,100 个 slot 每月约 $2,000,约相当于按需计费下处理 400 TB 的数据量。因此只有当月扫描量远超 200 TB 时,固定费率才划算。
slot 机制还影响并发表现:固定费率下若 100 个 slot 已被 50 个查询占满,第 51 个查询必须排队等待;而按需计费下 BigQuery 会根据查询需求动态分配更多 slot。另外,BigQuery 一般会缓存查询结果,模块演示中为得到一致结果通常会关闭缓存。
上手 BigQuery:查询公共数据集
BigQuery 内置大量可开箱即用的公共数据集。模块以纽约 Citi Bike 车站数据为例,在搜索栏中按表名找到citibike_stations(约 1,584 行、几 KB),即可用与big_query.sql第一段一致的 SQL 查询:
SELECT station_id, name FROM bigquery-public-data.new_york_citibike.citibike_stations LIMIT 100;查询结果会出现在界面底部的结果面板中,可导出为 CSV 或继续在 Data Studio 中探索。
外部表:让 BigQuery 直接读取 Cloud Storage 中的数据
模块使用纽约出租车行程数据(已上传至 Google Cloud Storage)。BigQuery 支持在 GCS 文件上创建外部表(external table):数据仍存放在 Cloud Storage,BigQuery 只保存元数据。
界面按「项目(project)/ 数据集(dataset)/ 表(table)」三级组织:taxi-rides-ny是项目、nytaxi是数据集、external_yellow_tripdata是表。创建外部表指向 2019、2020 两年的出租车 CSV:
CREATE OR REPLACE EXTERNAL TABLE `taxi-rides-ny.nytaxi.external_yellow_tripdata` OPTIONS ( format = 'CSV', uris = [ 'gs://nyc-tl-data/trip data/yellow_tripdata_2019-*.csv', 'gs://nyc-tl-data/trip data/yellow_tripdata_2020-*.csv' ] );创建成功后,BigQuery 会自动识别 CSV 的列名、类型与可空性,无需手工定义 schema(也可手动指定)。注意表详情中「长期存储 0 字节、表大小 0 字节、行数 0」——这是外部表的正常表现,因为数据在外部系统 GCS 中,BigQuery 无法直接获知行数与大小。
查询外部表与普通表无异:
SELECT * FROM taxi-rides-ny.nytaxi.external_yellow_tripdata LIMIT 10;返回结果包含 VendorID、上下车时间、乘客数、行程距离、上下车位置以及费用与总额等字段。
如何把数据送进 GCS
除了使用课程预置的 GCS 桶,仓库还提供了自行上传数据的脚本。03-data-warehouse/extras/目录下的 web_to_gcs.py 会从公开数据源下载各月出租车 CSV,用 pandas 按固定 dtypes 读取(保证 Parquet 列类型正确),转成 Parquet 后上传到 GCS:
- 在
extras目录执行uv sync安装依赖(见 pyproject.toml,包含 google-cloud-storage、pandas、pyarrow、python-dotenv、requests、tqdm); - 在
.env中设置GCP_GCS_BUCKET与GOOGLE_APPLICATION_CREDENTIALS(或使用 Google ADC); - 运行
uv run python web_to_gcs.py(简洁版)或uv run python web_to_gcs_with_progress_bar.py(带下载/转换/上传进度条)。
其中web_to_gcs_with_progress_bar.py的增强点值得注意:它会跳过 GCS 中已存在的对象、复用本地已下载的 CSV 与已转换的 Parquet,并通过分块读取(每块 10 万行)流式写入 Parquet,避免大文件内存溢出;上传时通过storage.blob._MAX_MULTIPART_SIZE与_DEFAULT_CHUNKSIZE调整分片大小以规避大文件超时。脚本默认执行web_to_gcs("2019", "green")、web_to_gcs("2020", "green")、web_to_gcs("2021", "green"),其中 2021 年 8 月起文件不存在属正常现象。
分区(Partitioning):按时间裁剪扫描量
分区是 BigQuery 最实用的优化手段之一。以 Stack Overflow 问题表(含创建日期、标题、标签等列)为例:若查询大多按日期过滤(如「只看 3 月的问题」),按创建日期分区后,每个日期成为一个独立分区。BigQuery 一旦确认只需读取某天的数据,就不会触碰其他分区的数据——处理的数据越少,成本越低。
回到出租车数据。先复制外部表,创建一张非分区表作为对照:
CREATE OR REPLACE TABLE taxi-rides-ny.nytaxi.yellow_tripdata_non_partitioned AS SELECT * FROM taxi-rides-ny.nytaxi.external_yellow_tripdata;复制数据需要一些时间(数据从 GCS 拷入 BigQuery 自有存储)。分区表只多一行PARTITION BY:
CREATE OR REPLACE TABLE taxi-rides-ny.nytaxi.yellow_tripdata_partitioned PARTITION BY DATE(tpep_pickup_datetime) AS SELECT * FROM taxi-rides-ny.nytaxi.external_yellow_tripdata;分区表能正确显示大小(约 13–14 GB),详情中可见其按tpep_pickup_datetime以「天」为粒度分区。一个快速判断表是否分区的技巧:在 schema 视图中,非分区表是一整块连续的列,分区表的列之间会出现一条分隔线。
现在对比两条等价的查询——统计 2019 年 6 月的 VendorID 去重结果:
SELECT DISTINCT(VendorID) FROM taxi-rides-ny.nytaxi.yellow_tripdata_non_partitioned WHERE DATE(tpep_pickup_datetime) BETWEEN '2019-06-01' AND '2019-06-30';非分区表预计处理1.6 GB(几乎是全表数据);将表名换成分区表,预计处理量骤降至约106 MB:
重复执行此类查询,每次只扫描 106 MB 而非 1.6 GB,成本直接随之下降。
还可以通过元数据检查各分区的行数分布。每个数据集都有INFORMATION_SCHEMA,其中的PARTITIONS视图可查询分区明细:
SELECT table_name, partition_id, total_rows FROM `nytaxi.INFORMATION_SCHEMA.PARTITIONS` WHERE table_name = 'yellow_tripdata_partitioned' ORDER BY total_rows DESC;该查询能看出哪些日期行数最多(示例表中 2019-02-01 最多),也可用于排查数据偏斜——某些分区数据量是否明显异常。
聚簇(Clustering):分区内的数据共置
聚簇解决的是分区内的数据组织问题。仍以 Stack Overflow 为例:表按日期分区后,再按标签聚簇,则同一分区内相同标签的行会物理相邻存储。由于相关行紧挨在一起,BigQuery 可以在分区内部跳过无关数据,进一步降低扫描量与查询耗时。
创建「既分区又聚簇」的表:
CREATE OR REPLACE TABLE taxi-rides-ny.nytaxi.yellow_tripdata_partitioned_clustered PARTITION BY DATE(tpep_pickup_datetime) CLUSTER BY VendorID AS SELECT * FROM taxi-rides-ny.nytaxi.external_yellow_tripdata;为何选VendorID作为聚簇列?因为该场景下查询总是同时按 vendor id 与上车日期过滤,分区与聚簇的组合正好命中这两个过滤条件。表详情页会确认:按tpep_pickup_datetime按天分区、按VendorID聚簇。
对比两种表在「统计 2019-06-01 至 2020-12-31 之间 Vendor 1 的行程数」上的表现:
SELECT count(*) as trips FROM taxi-rides-ny.nytaxi.yellow_tripdata_partitioned WHERE DATE(tpep_pickup_datetime) BETWEEN '2019-06-01' AND '2020-12-31' AND VendorID=1;运行前 BigQuery 对两张表的预估都是 1.1 GB——预估是近似值,聚簇能跳过多少数据只有运行后才知道。实际运行结果:仅分区的表处理 1.1 GB;分区 + 聚簇的表只处理约 843.5 MB。这就是聚簇的效果。记住这条经验:无论 BigQuery 运行前显示多少,都以运行后实际处理的字节数为准。
分区还是聚簇:取舍标准
分区与聚簇并非总是叠加使用,选择依据如下:
- 成本可预知性:分区的成本收益在查询前即可确定(按分区列过滤只读部分分区);聚簇的收益要到运行后才知道。BigQuery 允许为查询设置成本上限,超限则拒绝执行——只有成本可预知(即分区)时才能用此能力。
- 粒度:需要比分区更细的裁剪粒度时用聚簇。
- 管理能力:分区支持分区级管理(删除分区、在存储间迁移分区),聚簇没有。
- 列数:聚簇支持多列(常用多列过滤/聚合),分区只能基于单列。
- 基数(Cardinality):某列(或列组)去重值很多时适合聚簇——高基数是分区的障碍,因为单表分区数上限为4000。
分区选项与粒度
创建分区表时可选择分区依据(详见cohorts/2027/03-data-warehouse/02-partitioning-vs-clustering.md):
- 时间单位列(time-unit column):基于 timestamp/date 类型的列,如出租车上车时间;
- 摄取时间(ingestion time):按行写入时间分区,使用伪列
_PARTITIONTIME; - 整数范围(integer range):将整数列按区间切分。
时间单位列与摄取时间还支持粒度选择:日(默认)、小时、月、年。日粒度适合中等体量、数据在日期上均匀分布的场景;小时粒度适合海量数据按小时处理的场景(注意分区数上限 4000,可能需要过期策略清理旧分区);月/年粒度适合数据量小但日期跨度大的场景。
何时聚簇优于分区
以下情况应优先选择聚簇:
- 分区粒度会导致每个分区数据量过小(约小于 1 GB);
- 分区数会超过单表 4000 个的上限;
- 变更操作(mutation)频繁触及大部分分区(如每隔几分钟写一次)。
聚簇列的限制与自动重聚簇
聚簇最多指定4 列,必须是顶层、非重复(non-repeated)的列,可用类型为:DATE、BOOLEAN、GEOGRAPHY、INT64、NUMERIC、STRING、DATETIME。列顺序决定排序优先级:按 a、b、c 聚簇,则先按 a 排、再按 b、再按 c。
需要注意,分区与聚簇并非零成本:小于 1 GB 的小表两者都不会带来明显的性能提升,反而因元数据读取与维护增加开销,此时不设分区/聚簇更划算。此外,随着数据不断写入,新行的键范围可能与旧块重叠、削弱排序性质,BigQuery 会在后台自动执行自动重聚簇(automatic reclustering)恢复表的排序属性,且不影响查询性能、不额外计费;对分区表,聚簇在各分区范围内分别维护。
BigQuery 最佳实践清单
cohorts/2027/03-data-warehouse/03-bigquery-best-practices.md提供了一份围绕「降成本」与「提性能」两个目标的实践清单。
成本控制
- 避免
SELECT *:BigQuery 采用列式存储,每列独立存放。显式列出所需列时,BigQuery 只读这些列;用*则必须读取所有列。只取一两列时,这就是「几乎不读」与「读全表」的差别; - 运行前先看预估价格:查询编辑器右上角会显示预估费用,点击运行前先确认;
- 使用分区/聚簇表:让 BigQuery 跳过表的大部分数据;
- 谨慎使用流式插入(streaming inserts):会显著推高成本;
- 分阶段物化查询结果:如果一个 CTE(
WITH子句)在多处复用,不要每次重算,先跑一次把结果落成表,后续步骤直接引用该表。
另外,BigQuery 会缓存查询结果,重复执行相同查询时可直接命中缓存。
性能优化
- 始终在分区列或聚簇列上过滤,否则已建立的分区/聚簇形同虚设;
- 反规范化(denormalize)数据:数仓场景常与 OLTP 相反,把相关数据放在一起,避免查询时反复 join;
- 复杂结构用嵌套/重复列(nested/repeated columns):在避免彻底反规范化的同时保持相关数据共置;
- 合理使用外部数据源:从 GCS 等外部源读数据可能比读 BigQuery 自有存储更贵,不要过度使用;
- 先过滤再 JOIN:在 join 之前先缩小数据量;
- 不要把
WITH当作预编译语句:它只是单条查询内的命名子查询,不能跨查询复用; - 避免过度分片(oversharding):把数据拆成大量小表(如每天一张表)的效果劣于一张分区表。
更多提速技巧
- 避免使用 JavaScript 自定义函数(UDF);
- 用近似聚合函数替代精确函数,例如用 HyperLogLog++ 做去重计数;
ORDER BY放在查询最后;- 优化 join 顺序:把行数最多的表放第一位,行数最少的表次之,其余按体积递减排列——最大的表会被均匀分布到各节点,次大的表会被广播到所有节点。
BigQuery 内部原理:Colossus、Jupiter 与 Dremel
理解 BigQuery 的架构有助于设计自己的数据产品(详见cohorts/2027/03-data-warehouse/04-internals-of-bigquery.md)。三个核心组件:
- Colossus(存储):BigQuery 的底层分布式存储,采用列式格式,价格低廉。存储与计算分离是重大架构决策:数据增长只增加廉价的存储成本,昂贵的计算只在真正运行查询时发生;
- Jupiter(网络):连接存储与计算的数据中心网络,带宽约每秒 1 TB,使存储与计算分离后依然无延迟通信;
- Dremel(查询执行引擎):把查询拆成树状结构,逐层分发执行。
列式存储为什么重要
记录式(row-oriented)存储把整行作为一个整体(类似 CSV),简单易理解;列式(column-oriented)存储则让每一列独立存放。BigQuery 采用列式存储,带来的收益巨大:对列做聚合更快,且数仓查询本来就很少一次取全列——列分开存放后,其余列完全不被触碰,这也是SELECT *昂贵的原因。
Dremel 的查询执行树
以SELECT A, COUNT(B) FROM T GROUP BY A为例:root server接收查询后改写为SELECT A, SUM(C)(把计数变成对下层部分计数的求和),切分成切片分发给mixers(中间层);mixers 再细分给leaf nodes;leaf nodes 真正与 Colossus 交互、取数并在各自切片上执行操作,结果逐级汇总回 root server。正是这种分布式执行树,让 BigQuery 在数据规模增长时依然能线性扩展查询能力。
BigQuery ML:用 SQL 完成机器学习全流程
进阶主题之一是在 BigQuery 内部直接用 SQL 训练、评估、解释和调优机器学习模型(脚本见big_query_ml.sql,讲义见cohorts/2027/03-data-warehouse/05-machine-learning-in-bigquery.md)。本模块以「预测出租车小费金额」的线性回归模型为例。
为什么在 BigQuery 里做 ML
BigQuery ML 面向数据分析师与管理者:只需 SQL + 基础 ML 知识,无需掌握 Python/Java。另一个优势是数据不离仓——传统流程需要把数据导出到独立系统训练再部署,而 BigQuery 直接在仓库内建模型,省去导出步骤。
提示:机器学习对新手而言属于进阶内容,若暂时不熟悉可以跳过,不影响主流程。
定价
录制课程时的定价背景:数据存储 GB 免费;每月前 1 TB 查询免费;CREATE MODEL步骤每月前 10 GB 免费。免费额度之后,创建线性回归/逻辑回归/聚类/时间序列模型约 $250/TB,AutoML、DNN、boosted tree 模型约 $5/TB(另加 Vertex AI 训练费用)。以上为美国区价格,其他区域以官方为准。
特征选择与预处理
标签(要预测的目标)是tip_amount。特征从yellow_tripdata_partitioned表中挑选若干列:
SELECT passenger_count, trip_distance, PULocationID, DOLocationID, payment_type, fare_amount, tolls_amount, tip_amount FROM `taxi-rides-ny.nytaxi.yellow_tripdata_partitioned` WHERE fare_amount != 0;WHERE fare_amount != 0很关键:大量行程费用为 0,且这些行程小费几乎也为 0,会把模型带偏。
BigQuery ML 的预处理分自动与手动两种:自动预处理包括数值列标准化、类别列 one-hot 编码、数组 multi-hot 编码;手动预处理包括分桶(bucketization)、多项式展开、特征交叉、n-gram、min-max 缩放等。
关键陷阱:PULocationID、DOLocationID、payment_type虽然存储为整数,本质却是类别(位置 ID 264 不代表「264 倍」)。保持整数类型的话,BigQuery 会把它们当数值做标准化而不是 one-hot 编码。因此先建一张把这些列 CAST 成STRING的特征表:
CREATE OR REPLACE TABLE `taxi-rides-ny.nytaxi.yellow_tripdata_ml` ( `passenger_count` INTEGER, `trip_distance` FLOAT64, `PULocationID` STRING, `DOLocationID` STRING, `payment_type` STRING, `fare_amount` FLOAT64, `tolls_amount` FLOAT64, `tip_amount` FLOAT64 ) AS ( SELECT passenger_count, trip_distance, cast(PULocationID AS STRING), CAST(DOLocationID AS STRING), CAST(payment_type AS STRING), fare_amount, tolls_amount, tip_amount FROM `taxi-rides-ny.nytaxi.yellow_tripdata_partitioned` WHERE fare_amount != 0 );字符串列会被自动识别为类别特征并做 one-hot 编码。
创建模型
CREATE OR REPLACE MODEL `taxi-rides-ny.nytaxi.tip_model` OPTIONS (model_type='linear_reg', input_label_cols=['tip_amount'], DATA_SPLIT_METHOD='AUTO_SPLIT') AS SELECT * FROM `taxi-rides-ny.nytaxi.yellow_tripdata_ml` WHERE tip_amount IS NOT NULL;model_type='linear_reg'选择线性回归算法;input_label_cols指定标签列;DATA_SPLIT_METHOD='AUTO_SPLIT'让 BigQuery 自动划分训练集与评估集。训练约需 5 分钟,评估页会显示误差指标(示例中 MSE 约 8、MAE 约 1——模型不算最优,但本模块重在演示流程)。
特征检查、评估、预测与解释
BigQuery ML 提供四个配套函数:
ML.FEATURE_INFO:查看模型如何理解每个特征——数值列显示 min/max/mean(用于标准化),类别列显示类别数:
SELECT * FROM ML.FEATURE_INFO(MODEL `taxi-rides-ny.nytaxi.tip_model`);ML.EVALUATE:在指定数据集上打分,返回误差指标(示例中 MAE 约 1、MSE 约 150):
SELECT * FROM ML.EVALUATE(MODEL `taxi-rides-ny.nytaxi.tip_model`, (SELECT * FROM `taxi-rides-ny.nytaxi.yellow_tripdata_ml` WHERE tip_amount IS NOT NULL));ML.PREDICT:为数据集追加预测列predicted_tip_amount,可与真实tip_amount并排做人工评估:
SELECT * FROM ML.PREDICT(MODEL `taxi-rides-ny.nytaxi.tip_model`, (SELECT * FROM `taxi-rides-ny.nytaxi.yellow_tripdata_ml` WHERE tip_amount IS NOT NULL));ML.EXPLAIN_PREDICT:用STRUCT(3 as top_k_features)输出每条预测贡献最大的前 3 个特征(示例中三个类别特征——上下车位置与支付方式——对预测影响最大):
SELECT * FROM ML.EXPLAIN_PREDICT(MODEL `taxi-rides-ny.nytaxi.tip_model`, (SELECT * FROM `taxi-rides-ny.nytaxi.yellow_tripdata_ml` WHERE tip_amount IS NOT NULL), STRUCT(3 as top_k_features));超参数调优
模型不够理想时,标准手段是超参数调优。num_trials控制试验次数、max_parallel_trials控制并行试验数(加速整体运行)、l1_reg/l2_reg用hparam_range(连续范围)或hparam_candidates(候选列表)指定正则化参数搜索空间:
CREATE OR REPLACE MODEL `taxi-rides-ny.nytaxi.tip_hyperparam_model` OPTIONS (model_type='linear_reg', input_label_cols=['tip_amount'], DATA_SPLIT_METHOD='AUTO_SPLIT', num_trials=5, max_parallel_trials=2, l1_reg=hparam_range(0, 20), l2_reg=hparam_candidates([0, 0.1, 1, 10])) AS SELECT * FROM `taxi-rides-ny.nytaxi.yellow_tripdata_ml` WHERE tip_amount IS NOT NULL;CREATE MODEL还支持更多选项(学习率策略、early stopping、最小相对进步等),熟悉 ML 的读者可进一步深挖。
模型部署:导出 BigQuery ML 模型并用 Docker 提供服务
训练好的模型可以脱离 BigQuery 运行(步骤见extract_model.md,讲义见cohorts/2027/03-data-warehouse/06-deploying-a-machine-learning-model.md)。整体思路:用bq导出模型到 GCS → 用gsutil拷回本地 → 按 TensorFlow Serving 目录规范组织 → 用 Docker 启动服务 → 通过 HTTP 调用。
1. 导出模型到 Cloud Storage
gcloud auth login bq --project_id taxi-rides-ny extract -m nytaxi.tip_model gs://taxi_ml_model/tip_model导出完成后,tip_model目录会出现在taxi_ml_model桶中。
2. 拷贝到本地
mkdir /tmp/model gsutil cp -r gs://taxi_ml_model/tip_model /tmp/model导出的 BigQuery 模型本质是一个 TensorFlow 模型,包含assets、variables及若干元数据文件。
3. 用 TensorFlow Serving + Docker 提供服务
TensorFlow Serving 要求特定的目录布局:以模型名命名的目录,内含按版本号编号的子目录。因此创建serving_dir/tip_model/1并把模型文件拷入版本目录:
mkdir -p serving_dir/tip_model/1 cp -r /tmp/model/tip_model/* serving_dir/tip_model/1 docker pull tensorflow/serving docker run -p 8501:8501 \ --mount type=bind,source=`pwd`/serving_dir/tip_model,target=/models/tip_model \ -e MODEL_NAME=tip_model -t tensorflow/serving &参数含义:-p 8501:8501把容器 8501 端口映射到本机;--mount把本地 serving 目录绑定挂载进容器;-e MODEL_NAME=tip_model告诉 TensorFlow Serving 服务哪个模型。
4. 检查模型状态
TensorFlow Serving 暴露 REST API。向http://localhost:8501/v1/models/tip_model发 GET 请求(视频中使用 Postman),响应应显示 tip_model 版本 1 状态为 AVAILABLE。
5. 通过 HTTP 预测
向http://localhost:8501/v1/models/tip_model:predictPOST 一条 JSON,字段与训练特征一致:
curl -d '{"instances": [{"passenger_count":1, "trip_distance":12.2, "PULocationID":"193", "DOLocationID":"264", "payment_type":"2","fare_amount":20.4,"tolls_amount":0.0}]}' \ -X POST http://localhost:8501/v1/models/tip_model:predict示例中该行程预测小费约 $3.2;把payment_type改成 2 后预测值骤降至约 $0.26。至此完成闭环:BigQuery 内用 SQL 训练 → 导出 → Docker 容器中以 REST 服务对外提供预测。
模块作业:把 2024 年上半年出租车数据搬进 BigQuery
README 将作业指向 2026 年模块三作业:使用2024 年 1 月至 6 月的 Yellow Taxi 行程数据(Parquet 格式),完成「建外部表 → 建非分区/分区聚簇表 → 对比查询扫描量」的练习。注意:本次作业创建外部表时必须使用 PARQUET 格式选项。
数据加载可使用配套的 load_yellow_taxi_data.py 脚本:它通过google.cloud.storage客户端、以 4 线程并发从公开 CDN 下载 6 个月份的 Parquet 文件(每个 8 MB 分块),自动创建(或校验)GCS 桶,上传后调用storage.Blob(...).exists()做存在性校验,失败自动重试 3 次(间隔 5 秒)。使用前需准备带GCS Admin权限的 Service Account(或通过 gcloud SDK 认证),并修改脚本中的BUCKET_NAME。
作业配套 SQL 见big_query_hw.sql:用 FHV 数据演示了外部表创建、count(*)全表统计、COUNT(DISTINCT(dispatching_base_num))去重统计、非分区表与「按dropoff_datetime分区 + 按dispatching_base_num聚簇」表的对比查询。
延伸资源
- 模块幻灯片链接及全部课程视频入口见 03-data-warehouse/README.md;
- 基础 SQL 脚本:
big_query.sql; - ML 脚本:
big_query_ml.sql; - 作业脚本:
big_query_hw.sql; - 模型导出与部署手册:
extract_model.md; - 数据上传工具:
extras/web_to_gcs.py与extras/web_to_gcs_with_progress_bar.py; - 逐讲文字讲义:
cohorts/2027/03-data-warehouse/下01-data-warehouse-and-bigquery.md至06-deploying-a-machine-learning-model.md; - README 的 Community notes 区块还汇集了历届学员的社区笔记与 2024 年视频转录稿链接,可作为学习参考。
【免费下载链接】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),仅供参考