18 大数据的处理与学习

本章会用到的数据

  • 公开数据百万行教学交易数据用于分块计算;恒瑞医药分钟行情用于真实数据练习。

  • 这些数据能做什么:比较 Pandas 与 Dask 的读取方式,并学习在数据较大时分块计算。

  • 分析时注意:先只读取需要的列,观察运行时间和内存;时间预测仍要按先后顺序划分样本。

  • 示例说明:41 MB 的教学文件用于理解方法,不代表真正的超内存计算。

【课堂核心】3 学时学习安排(1/2)

学习内容 分钟
动机与先修 20
概念与公式 35
输入说明与检查 20

【课堂核心】3 学时学习安排(2/2)

学习内容 分钟
独立任务与反馈 45
完整答案与完整示例 40
本章小结 20

【可选拓展】:集群实际使用、调度器细节与大规模机器学习框架比较。

欢迎来到大数据时代

  • 我们正处在一个数据爆炸的时代。
  • 从金融交易到社交媒体,海量数据正在重塑商业决策。
  • 但我们迄今为止所学的工具,在面对“真正的大数据”时,会遇到一堵看不见的墙。
  • 今天,我们的任务就是——翻越这堵墙

本章学习议程

我们将循序渐进,探索大数据的世界:

  1. 问题的根源:为什么传统方法会失效?
  2. 核心思想:解决大数据问题的通用策略。
  3. 关键工具(Dask):深入了解Python生态的并行计算利器。
  4. 动手实践:用 Dask 对公开真实分钟数据执行多分区惰性计算。
  5. 高级话题:区分教学缩放接口、真实分区计算与需要另行资源验证的集群机器学习。

学习目标 (1): 理解挑战

在这堂课结束时,我希望你们能够理解

为什么处理海量数据是一个巨大的挑战,以及它与传统数据分析的根本区别。

传统内存流程面对海量数据的瓶颈完整数据拼图涌入容量有限的漏斗,只有一个脱离上下文的碎片通过,右侧问号表示输出残缺而无法支持分析。 海量数据 传统分析方法 内存限制 ? 无用 / 残缺的信息

学习目标 (2): 了解工具

并且,能够了解

Dask和Spark这两个业界主流的大数据处理框架的核心思想。

Dask与Spark的分解并行思想一个大任务被拆成多个分区,分发到并行工作节点处理,再汇总为结果;图示强调分解、调度和聚合而非框架性能结论。 大任务 Dask Spark 并行处理结果

学习目标 (3): 掌握实践边界

最终,能够实践掌握

  • 使用Dask执行基本的数据操作,并体会其与Pandas的异同。
  • 解释海量数据机器学习的分区、增量与分布式训练策略,并在单机多分区示例中核对接口和执行边界。
掌握大数据实践 一个代码块被执行,将混乱的数据点转化为清晰的图表结果。 > import dask.dataframe as dd > ddf = dd.read_csv('big.csv') > result = ddf.mean().compute() 清晰的结果

问题的根源

核心问题:当数据大到内存装不下时,该怎么办?

我们之前所有的数据分析工作,都有一个隐含的前提…

那个我们从未质疑过的前提…

所有数据都能一次性加载到计算机的内存(RAM)中。

pandas 就是这样工作的。当你调用 pd.read_csv() 时,它会尝试将整个文件读入内存。

这在小数据集上表现优异,但在大数据时代,这个前提本身就是问题的根源。

内存之墙:理想与现实

我们希望用我们有限的内存,去分析无限增长的数据。

内存与数据的矛盾 一个小盒子代表的RAM面对着一个巨大的、无边无际的数据流,中间被一堵墙隔开。 RAM 有限 DATA 无限增长 内存之墙

个人电脑的内存是极其有限的

让我们来看一下现实:

  • 学生笔记本电脑:通常是 8GB16GB 内存。
  • 高端工作站:也许能达到 64GB 或 128GB,但很少超过 200GB。

这就像…

一个形象的比喻

…试图用一个水杯去装整个游泳池的水。

水杯与游泳池 一个小水杯面对一个巨大的游泳池,比喻内存和数据的关系。 16GB RAM TB级数据

假设情景:十亿用户级交易数据

下面只做容量级别的课堂估算,不对应任何具名平台的披露数据。四项输入均为假设,用途是判断单机内存约束,而不是陈述真实业务规模。

  • 假设月活跃用户: \(10^9\)
  • 假设每位用户每月交易: 50 次
  • 假设每次交易的数值字段: 10 个
  • 假设每个字段存储大小: 8 字节(float64)

这会产生多大的数据量?

触目惊心的数据量估算

让我们来做个简单的计算:

\[ \large{ \text{数据量} = \underbrace{10^9}_{\text{用户数}} \times \underbrace{50}_{\text{交易数}} \times \underbrace{10}_{\text{变量数}} \times \underbrace{8}_{\text{字节/变量}} } \]

问题的严重性:4TB

计算结果是:

\[ \large{ \text{结果} \approx 4 \times 10^{12} \text{ 字节} \approx \textbf{4 Terabytes} } \]

4TB 的裸数值容量通常远超学生常用设备的内存,且还未计入索引、字符串和运行时开销。

可选策略包括按列/分区读取、数据库下推、流式聚合和分布式计算;是否需要 Dask 取决于数据布局、操作类型与硬件预算。

解决方案的核心思想

解决方案:分而治之 (Divide and Conquer)

当数据无法一次性装入可用内存时,一个常见办法是:

将一个大任务(处理整个数据集)分解成许多可以在小块数据上执行的小任务,然后将结果汇总起来。

这正是大数据处理框架的精髓所在。

“分而治之”的可视化理解

分而治之 一个大的数据块被分解成小块,处理后结果再合并。 大数据 1. 分解 2. 并行处理 3. 汇总 最终结果

Python生态中的两大主流工具

在Python世界里,有两个强大的程序库专门用来实现“分而治之”的思想:

  1. Dask: 一个轻量级、原生支持Python的并行计算库。
  2. Apache Spark: 一个功能更强大、生态更成熟的分布式计算引擎。

我们今天将主要以Dask为例进行讲解,因为它对熟悉Pandas的你来说最容易上手。

工具介绍: Dask

Dask Logo Dask的蛇形标志,象征其Python原生特性。

Dask: 轻量级、Python原生、与现有库(Pandas, NumPy)无缝集成,学习曲线平缓。

工具介绍: Apache Spark

Apache Spark Logo Spark的星火标志,象征其强大的计算能力。

Apache Spark: 行业标准、功能强大、生态成熟、支持多语言(通过PySpark在Python中使用)。

深入Dask: 为Python设计的并行计算库

Dask: 熟悉的味道,更强的能力

Dask是一个开源项目,旨在将Python数据科学生态系统(Pandas, Scikit-Learn, NumPy)的能力扩展到大数据和并行计算领域。

它的设计哲学是:尽可能地模仿现有库的API,让你用最小的学习成本处理更大的数据集。

Dask的优点:与Pandas的惊人相似性

Dask最吸引人的一点在于,你不需要学习一套全新的语法。

  • Dask DataFrame: Dask的核心数据结构之一,它在内部由多个小的Pandas DataFrame组成。
  • API兼容: Dask DataFrame 覆盖许多常用 Pandas 操作,但并非完整替代;不受支持的操作、全局排序与大规模 shuffle 需要改写或换用数据库/分布式方案。

这使得从Pandas迁移到Dask变得异常平滑。

Dask的五大核心特点

  1. 懒惰执行 (Lazy Execution)
  2. 任务调度 (Task Scheduling)
  3. 智能内存管理 (Smart Memory Management)
  4. 并行与分布式计算 (Parallel & Distributed)
  5. 无缝生态集成 (Ecosystem Integration)

我们将逐一深入探讨这些概念。

核心特点 (1): 懒惰执行 (Lazy Execution)

这是Dask与Pandas最本质的区别。

  • Pandas: 你写的每一行代码都会立即执行
  • Dask: 多数 collection 的读取与变换先构建任务图,不会在定义表达式时立刻物化全部结果。

相反,Dask会构建一个任务图(Task Graph),记录下你想要执行的所有计算步骤。

什么是任务图?

任务图是 Dask 的“作战计划”:

  • 需要执行的所有任务(例如:从文件中读取一个数据块,过滤行,计算平均值)。

  • 任务之间的依赖关系(例如:必须先读取数据块,然后才能对其进行过滤)。

  • 触发执行.compute().persist().head() 等查看或缓存操作。

  • 模型触发:需要统计量归约的 estimator .fit() 也会执行任务图。

  • 判断原则:逐项识别 API 的执行与物化边界,不能简化为“只有 .compute() 才计算”。

思想实验:懒惰执行就像是写一份菜谱

想象一下做一道复杂的菜:

立即执行 (Pandas)

  • 读第一步“切洋葱”
  • 立刻切洋葱
  • 读第二步“热锅”
  • 立刻热锅
  • …一步一动…

懒惰执行 (Dask)

  • 读完整份菜谱
  • 🧠 在脑中规划
  • “先切好所有蔬菜,再混合酱料,最后一起下锅。”
  • 直到决定“开火”时,才真正开始动手。

Dask的这种方式让它有机会在执行前对整个计算流程进行优化。

核心特点 (2): 任务调度 (Task Scheduling)

懒惰执行构建了任务图,而任务调度器(Task Scheduler)则是负责执行这个图的“大脑”。

Dask会将大的计算任务分解成许多小的、独立的任务,然后智能地安排它们的执行顺序,并分配到不同的CPU核心或机器上并行处理。

任务调度器的工作原理

Dask任务调度器 调度器将任务图中的任务分配给多个工作核心。 调度器 任务图 (Plan) CPU 核心 1 CPU 核心 2 CPU 核心 ...N 分配任务

核心特点 (3): 智能内存管理

Dask可按分区调度任务,但内存行为取决于图依赖、分区大小、缓存与调度器配置。

计算时相关分区会进入内存;不再被依赖的中间结果可能被释放,分布式调度器也可按阈值 spill 到磁盘,但持久化对象、过大分区或 shuffle 仍可能造成内存压力。

因此必须监控峰值内存并合理划分分区,不能假定每个数据块计算后都会立即释放。

内存管理的“流水线”模式

Dask内存管理 数据块按任务依赖进入内存;无后续依赖的中间结果可能释放,也可能因缓存或调度策略保留。 磁盘 内存 (RAM) 结果 处理中...

核心特点 (4): 并行与分布式计算

Dask的强大之处在于它的可伸缩性:

并行计算 (Parallel)

  • 单机多核是可选执行方式
  • 只有任务可并行、分区足够大且调度、序列化与 I/O 开销相对较小时,合适的 Dask scheduler 才可能带来加速;Python GIL、任务过碎或数据搬运都可能抵消收益。

分布式计算 (Distributed)

  • 多机集群扩展容量与并发
  • Dask 可以协调多台机器处理分区任务,但收益必须用同一工作负载的 wall time、峰值内存、吞吐与调度开销实测;数据规模大不保证更快。

并行 vs. 分布式

并行计算与分布式计算 左侧显示单台机器内的多核并行,右侧显示多台机器的分布式集群。 并行计算 (单机) Core1 Core2 Core3 ... ... CoreN 分布式计算 (集群) 机器1 机器2 ...N 网络

核心特点 (5): 与现有生态系统无缝集成

Dask不仅仅是模仿了Pandas,它与整个Python科学计算栈都深度集成。

这意味着你可以用很小的代码改动,就将现有分析流程从处理小数据升级为处理大数据。

Dask 与 Python科学计算全家桶

Dask生态系统集成 Dask作为中心,与Pandas, NumPy, Scikit-Learn等库协同工作。 Dask Pandas NumPy Scikit Learn ...and more

Dask 与 Pandas 的操作对比

下表直观地展示了Dask与Pandas在基本操作上的相似与不同之处。

操作 Pandas Dask
数据结构 DataFrame, Series Dask DataFrame, Dask Series
导入库 import pandas as pd import dask.dataframe as dd
读取数据 pd.read_csv(...) dd.read_csv(...)
查看数据 df.head(), df.tail() df.head(), df.tail()
数据过滤 df[df['col'] > 0] df[df['col'] > 0]
分组聚合 df.groupby('col').sum() df.groupby('col').sum()
数据合并 pd.concat(), df.merge() dd.concat(), df.merge()
计算执行 立即执行 多数变换先建任务图.compute().persist().head() 或 estimator .fit() 等操作可触发执行

理解.compute(): Dask操作的“执行”按钮

.compute()方法是Dask懒惰执行模型的关键。

  • 当你写下 df_filtered = df[df['column'] > 0] 时,如果df是一个Dask DataFrame,df_filtered不是一个包含结果的数据帧。
  • df_filtered此时只是一个“任务图”,一个“计算计划”,它描述了如何得到过滤后的结果。
  • result = df_filtered.compute() 触发本例的读取、过滤与汇总,并返回 Pandas DataFrame
  • 其他流程也可能由 .persist()、查看操作或 estimator .fit() 触发。

.compute() 的作用:从计划到现实

.compute()的作用 Dask对象通过.compute()方法转换成Pandas对象。 Dask DataFrame (计算计划) Pandas DataFrame (真实结果) .compute()

让我们进入实践:用Dask分析大型数据集

实践环节目标

  1. 安装 Dask 库。
  2. 检查仓库随附的教学缩放交易数据版本,并明确它不代表真实市场或超内存规模。
  3. 使用Dask读取并分析该数据集。
  4. 亲身体验懒惰执行.compute()

步骤1: 安装 Dask 库

打开你的终端或Anaconda Prompt,根据你的需求选择安装命令:

  • 安装完整版 (推荐): 包含所有依赖项,如分布式调度、可视化等。

    pip install "dask[complete]"
  • 仅安装核心库:

    pip install dask

步骤2: 检查确定的教学缩放数据集

本章使用约 41MB 的公开教学 CSV,观察分区读取、惰性执行和分块计算。它通常小于机器内存,因此本章只练习方法,不声称完成了 TB 级数据测试。

我们的目标是对同一份一百万行数据版本完成“输入检查—分区读取—年度归约—期望输出”的可自查小结。

开始前先检查数据文件

代码
# 导入 `pandas`,用于查看字段、行数与时间范围。
import pandas as pd
# 为后续 Dask-ML 与公开分钟验证导入 `numpy` 并绑定 `np`,用于数组和数值计算。
import numpy as np
# 导入 `os`,用于检查文件是否存在与查看文件大小。
import os
from pathlib import Path
from urllib.request import urlretrieve  # 复用本章隐藏设置单元安装的浏览器标识下载器

# 固定教学 CSV 的期望行数,作为截断或追加文件的结果不理想时如何解释。
n_rows = 1_000_000 # 100万行
# 文件不存在时,从与其他章节相同的公开数据域下载。
# 按 Linux 共享数据、Windows 共享数据、项目缓存的顺序选择大型交易教学文件。
file_path = next((candidate_path for candidate_path in [Path('/home/ubuntu/r2_data_mount/data/course/large_transactions.csv'), Path('C:/qiufei/data/course/large_transactions.csv'), Path('data/course/large_transactions.csv')] if candidate_path.exists()), Path('data/course/large_transactions.csv'))
if not file_path.exists():
    file_path.parent.mkdir(parents=True, exist_ok=True)
    urlretrieve('https://assets.qiufei.site/data/course/large_transactions.csv', file_path)

查看数据的字段与范围

先查看字段、行数与时间范围,确认下载的数据适合后续练习。

代码
# 读取字段名而不载入数据行,核对教学 CSV 的五列说明与顺序。
observed_columns = pd.read_csv(file_path, nrows=0).columns.tolist()
# 要求字段说明完全一致,防止旧版或手工修改文件进入 Dask 练习。
assert observed_columns == ['transaction_id', 'timestamp', 'stock_id', 'price', 'volume'], '教学 CSV 字段或顺序不符合课程说明'
# 初始化逐块累计的行数与时间上下界,避免为检查一次性物化全部字段。
validated_row_count, validated_time_min, validated_time_max = 0, None, None
# 分块读取交易编号与时间字段,检查完整样本范围而不依赖文件字节数。
for validation_chunk in pd.read_csv(file_path, usecols=['transaction_id', 'timestamp'], chunksize=100_000):
    # 把当前块时间解析为时间类型,供全文件上下界比较。
    validation_times = pd.to_datetime(validation_chunk['timestamp'])
    # 累计当前块行数,最终应严格等于一百万行。
    validated_row_count += len(validation_chunk)
    # 更新已读分块中的最早时间,核对数据版本说明的确定时间窗口。
    validated_time_min = validation_times.min() if validated_time_min is None else min(validated_time_min, validation_times.min())
    # 更新已读分块中的最晚时间,核对数据版本说明的确定时间窗口。
    validated_time_max = validation_times.max() if validated_time_max is None else max(validated_time_max, validation_times.max())
# 要求输入严格包含一百万行,防止截断文件或追加文件通过检查。
assert validated_row_count == n_rows, '教学 CSV 行数不符合课程说明'
# 要求所有时间均落在确定的左闭右开秒级窗口内。
assert validated_time_min >= pd.Timestamp('2017-07-14 02:40:00') and validated_time_max < pd.Timestamp('2020-09-13 12:26:40'), '教学 CSV 时间范围不符合课程说明'
# 显示数据概况,随后进入 Dask 读取与练习。
print({'字段数': len(observed_columns), '行数': validated_row_count, '时间范围': (validated_time_min, validated_time_max)})
{'字段数': 5, '行数': 1000000, '时间范围': (Timestamp('2017-07-14 02:42:24'), Timestamp('2020-09-13 12:25:54'))}
代码
# 读取教学 CSV 的实际字节数并换算为 MB,核对 8MB 分区参数会产生多个块。
file_size_mb = os.path.getsize(file_path) / (1024**2)
# 报告这份 CSV 占用的兆字节数,核对 8MB 分区设置能产生多个分区。
print(f'文件大小约为: {file_size_mb:.2f} MB')
文件大小约为: 40.44 MB

为什么先看数据概况?

  • 字段、行数与时间范围能帮助我们发现下载不完整或文件用错等问题。
  • 分块检查只让一个十万行检查块同时驻留,演示降低检查阶段峰值对象的模式。
  • 本次约 41MB 文件没有证明一次性读取会 OOM;是否 OOM 必须结合实测峰值和机器预算。
  • 后续年度 groupby 仍需汇总所有源分区,最终四行输出不等于只读取少数分区。

步骤3: 使用 Dask 读取数据

现在,关键时刻来了。我们不用pandas,而是用dask.dataframe来读取这个大文件。

代码
# 为“步骤3: 使用 Dask 读取数据”导入 `dask.dataframe` 并绑定 `dd`,用于建立当前任务的分区表格与惰性任务图。
import dask.dataframe as dd

# 使用Dask读取CSV
# 8MB分块可保证本教学文件形成多个分区
ddf = dd.read_csv('large_transactions.csv', blocksize='8MB')

# 此时尚未物化全量分区;`read_csv` 可能已抽取少量样本推断字段和类型。
print('Dask DataFrame已创建!')
# 打印 Dask DataFrame 的逻辑分区数,核对 8MB 块设置确实形成多分区教学输入。
print(f'分区数: {ddf.npartitions}')
# 要求教学文件确实跨越多个 Dask 分区,否则不能据此演示分块执行。
assert ddf.npartitions > 1, '教学文件必须形成多个分区才能演示分块执行'
Dask DataFrame已创建!
分区数: 5

解读 dd.read_csv

  • import dask.dataframe as dd: 这是Dask DataFrame模块的惯用别名,就像import pandas as pd一样。
  • blocksize='8MB':该教学参数确保约 25—45MB 的文件形成多个分区。实际使用时的参数应依据文件大小、可用内存和调度开销实测选择。

审视Dask DataFrame对象

让我们看看刚刚创建的ddf对象到底是什么。

代码
# 输出惰性 Dask DataFrame 元数据,查看分区、字段与尚未物化的任务对象。
print(ddf)
Dask DataFrame Structure:
              transaction_id timestamp stock_id    price volume
npartitions=5                                                  
                       int64    string    int64  float64  int64
                         ...       ...      ...      ...    ...
...                      ...       ...      ...      ...    ...
                         ...       ...      ...      ...    ...
                         ...       ...      ...      ...    ...
Dask Name: to_string_dtype, 2 expressions
Expr=ArrowStringConversion(frame=FromMapProjectable(b606a15))

解读输出结果

  • 没有数据! 和Pandas不同,它不会打印出数据内容。
  • 列名和类型: 它显示了Dask DataFrame Structure,包含了列名和dtype。Dask通过智能地读取文件的一小部分来推断这些元信息。
  • 分区 (npartitions): 这是核心。它告诉你这个Dask DataFrame在逻辑上被分成了多少块(分区)。这个数量取决于文件总大小和blocksize

逻辑依赖示意:从分区读取到聚合

Dask 先构造惰性任务计划,再由 .compute() 触发本例计算。下面的图只用 3 个代表性分区说明过程;当前 CSV 实际形成 5 个分区。

代码
ddf['amount'] = ddf['price'] * ddf['volume']  # 为每个分区建立“价格乘数量”的惰性列变换
total_amount = ddf['amount'].sum()  # 建立跨分区归约任务,此时尚未执行全量计算
print({'当前实际分区数': ddf.npartitions, '归约对象类型': type(total_amount).__name__})  # 检查五个实际分区与等待执行的 Dask 标量,不把简化 SVG 冒充完整任务图
{'当前实际分区数': 5, '归约对象类型': 'Scalar'}
Dask 分区变换与归约的简化逻辑依赖示意 左列三个圆角矩形代表实际五个分区中的三个代表性读取任务,中列圆角矩形代表逐分区计算价格乘数量,右列圆角矩形代表跨全部分区求和;箭头从左向右表示依赖与汇聚。此图不是完整 Dask graph 检查输出。 读取分区 1读取分区 2读取分区 … N price × volumeprice × volumeprice × volume 跨分区求和compute() 触发

如何解读逻辑依赖示意

本图用三个代表性分区解释逻辑依赖;当前对象实际有 5 个分区,完整任务键、优化层和调度细节未画出。

  • 左列圆角矩形:读取源分区;… N 表示其余分区采用同一逻辑。
  • 中列圆角矩形:在各分区内计算 price × volume,仍只形成惰性列变换。
  • 右列圆角矩形:汇总全部分区的局部结果,形成跨分区归约。
  • 从左向右的箭头:表示读取 → 局部变换 → 全局汇聚的依赖方向。

构建 amountsum 只形成计划;本例调用 .compute() 才触发该归约。其他查看操作、.persist() 或 estimator .fit() 也可能在各自流程中触发执行。

步骤4: 执行计算并获取结果

现在,让我们通过调用.compute()来触发这些计算。我们将计算price列的基本描述性统计。

代码
# 计算描述性统计
# 这个过程会花费一些时间,因为Dask现在正在读取所有数据块并进行计算
print("开始计算描述性统计...")
# 由 `.compute()` 触发全部价格分区的描述统计归约,并返回物化的 Pandas Series。
summary_stats = ddf['price'].describe().compute()
# 报告聚合结果已经物化为 Pandas 对象。
print("计算完成!")
# 输出物化后的价格描述统计,检查跨分区归约的最终 Pandas 结果。
print(summary_stats)
开始计算描述性统计...
计算完成!
count    1000000.000000
mean         255.179348
std          141.401938
min           10.000000
25%          133.150000
50%          255.700000
75%          377.970000
max          500.000000
Name: price, dtype: float64

.compute() 幕后发生的事

  1. Dask Scheduler接收到 describe 任务图。
  2. 它将任务(为每个分区计算统计量)分配给可用的CPU核心。
  3. 每个核心读取一个数据块,计算 count, mean, std, min, max 等。
  4. 所有分区的中间结果被汇总起来,计算出最终的全局统计量。
  5. 峰值内存取决于分区大小、聚合中间量、任务图、调度器与 spill 配置;必须实测,不能由“分区执行”推断严格的小内存上界。

Dask的Pandas-like操作:分组聚合

Dask的groupby操作也和Pandas几乎一样。让我们计算每只股票的平均交易量。

代码
# 计算每只股票的平均交易量
# 同样,这只是定义了计算计划
mean_volume_by_stock = ddf.groupby('stock_id')['volume'].mean()

# 使用查看操作把输出限制为前十个分组结果,以节省屏幕空间。
# `head(10)` 返回 Pandas 对象,但全局 groupby 均值仍可能需要所有源分区提供分子与分母。
print("开始分组聚合计算...")
# 取得前十个已聚合分组;`head` 限制返回行与输出分区,不承诺只读取前几个源分区。
result = mean_volume_by_stock.head(10)
# 报告受限输出已由 `head` 触发并物化;全局归约的源分区工作量不能由返回十行判断。
print("计算完成!")
# 输出前十个股票的平均成交量,检查分组键与聚合值。
print(result)
开始分组聚合计算...
计算完成!
stock_id
1     4951.722832
2     5243.609553
3     5261.428571
4     5137.991812
5     5016.143145
6     5053.174588
7     4856.658436
8     5126.087500
9     5171.106833
10    5086.707676
Name: volume, dtype: float64

简单小结:Dask vs. Pandas

  • 相似性: API非常接近,如果你会Pandas,你几乎就会用Dask DataFrame。
  • 核心差异: Dask 的多数 collection 变换先构建任务图,再由 .compute().persist()、部分查看操作或 estimator .fit() 等触发执行;分区与调度使其能够处理大于单机内存的数据。
  • 适用场景: Dask 常适合 Python 原生、可分区的单机或集群任务;是否采用它必须同时核对数据格式、操作与 shuffle 图、峰值内存、容错、集群条件、团队生态和基准测试,文件大小本身不是决定规则。

超越单机:Apache Spark简介

简介:Apache Spark 与 PySpark

当任务需要成熟的多语言数据平台、强容错与大规模集群治理时,可以评估 Apache Spark;数据规模只是输入之一,不能仅凭 TB 数量替代对格式、操作图、倾斜、资源与团队生态的判断。

  • Apache Spark: 是一个为大规模数据处理而设计的统一分析引擎。它最初用Scala编写,拥有卓越的性能和可伸缩性。
  • PySpark: 是Spark为Python提供的官方API,让我们可以在Python中利用Spark的强大功能。

Spark的核心特点:内存集群计算

Spark最著名的特点是其内存计算(In-memory Computing)能力。

与早期的大数据系统(如Hadoop MapReduce)不同,Spark可以将中间计算结果保存在集群中各台机器的内存里,大大加快了迭代算法(如机器学习)和交互式查询的速度。

Spark的集群计算模型

Spark集群计算模型 一个主节点协调多个工作节点,每个节点都有自己的CPU和内存。 主节点 (Driver) 工作节点 1 工作节点 2 工作节点 N

Dask 与 Spark: 何时选择哪个?

这是一个常见的问题。下表先比较两种工具;实际选择还要结合数据格式、运行时间、内存和已有学习环境,不能只看文件大小。

特性 Dask Apache Spark
生态系统 纯Python,与Pandas/Numpy/Scikit-Learn紧密集成 多语言 (Scala, Java, Python, R),更庞大、更复杂的生态
设置复杂度 非常简单 (pip install dask) 相对复杂,通常需要配置Java环境和集群管理器
适用场景 单机并行、中等规模集群,Python原生项目 大型企业级集群,超大规模数据,多语言团队
学习曲线 平缓,对Pandas用户友好 较陡峭,需要理解Spark的架构和概念

总结:如何选择

统一决策证据

  • 先记录数据格式与分区键、操作和 shuffle 图、倾斜风险、实测峰值内存与墙钟时间。
  • 再比较失败恢复与容错要求、可用集群、既有 Python/Spark 生态和团队运维能力。
  • Dask 与 PySpark 都必须通过代表性切片和目标规模压测;数据量不单独决定框架。

大数据时代的机器学习

大数据机器学习的挑战

  • 失败条件:estimator 要求全量物化,且输入或关键中间量超过内存预算。
  • 典型后果.fit(X, y) 发生 OOM 或严重换页。
  • 可扩展接口partial_fit、核外 solver,或可接受 Dask 分区输入的 estimator。
  • 仍需核对:具体 solver、训练规则与中间量规模。

何时全量 .fit(X, y) 可能失败?

模型拟合失败 当估计器要求物化全量输入且数据与中间量超过内存预算时,巨大数据块无法装入较小内存并可能溢出。 X_train (500GB) RAM (16GB) model.fit() MemoryError!

对策1 (简单但有缺陷): 随机抽样

最直接的方法是将大数据集缩小到内存可以容纳的大小。

  • 做法: 从海量数据中随机抽取一部分样本(例如1%)。
  • 优点: 简单易行,可以使用所有标准的scikit-learn工具。

随机抽样的代价

缺点:

  1. 信息丢失: 丢弃了99%的数据,可能会错过重要的模式,尤其是在处理稀有事件(如金融欺诈)时。
  2. 模型性能: 对于需要大量数据才能学习的复杂模型(如深度学习),抽样后的数据可能不足以训练一个高性能的模型。
随机抽样的信息丢失 大部分数据被丢弃,只有一个小样本被保留。 全部数据 99% 被丢弃! 保留 1% 样本

对策2:一种常见的同步数据并行训练

下面展示的是同步数据并行的一种常见设计,而不是分布式训练的完整定义:各 worker 持有模型副本、处理不同批次,并在每步聚合梯度后同步参数。模型并行、流水线并行与异步参数服务器等设计不遵循同一流程。

它可利用多个计算单元扩展训练容量或缩短训练时间,但实际收益取决于计算/通信比、互连带宽、负载均衡、同步等待和容错开销,必须以单机基线实测。

同步数据并行示例 (1/4):分区与分发

  1. 准备环境: 配置一个包含多个计算节点(worker)的集群。
  2. 数据分区: 将大型数据集分割成许多小批次(batches)。
  3. 数据分配: 将这些数据批次均匀地分配给每个节点。
数据分区与分发 大数据被切分并发送到不同的工作节点。 大数据 Worker 1 Worker 2 Worker 3

同步数据并行示例 (2/4):模型复制

  1. 模型复制: 在每个节点上初始化一个完全相同的模型副本(拥有相同的初始权重)。
模型复制 一个模型模板被复制到所有工作节点。 模型

同步数据并行示例 (3/4):并行计算

  1. 并行计算: 每个节点使用自己分配到的数据批次,独立地对模型进行训练,并计算出模型参数的更新量(梯度)。
并行计算梯度 每个工作节点独立计算梯度。 + 计算 梯度 G1 + 计算 梯度 G2 + 计算 梯度 G3

同步数据并行示例 (4/4):聚合与更新

  1. 梯度聚合: 所有节点计算出的梯度被汇总起来。
  2. 模型更新: 使用聚合后的梯度来更新所有节点上的模型参数。
  3. 迭代: 重复并行计算和更新的步骤。
聚合与更新 所有梯度被发送到中心进行聚合,然后更新所有模型。 聚合梯度 已更新

分布式训练的优化与工具:Dask-ML

幸运的是,我们有现成的工具来处理这些复杂性。

Dask-ML是 Dask 生态中扩展机器学习步骤的模块,部分接口接近 scikit-learn;是否真正形成分布式或核外训练,仍取决于 estimator、solver、collection、scheduler 与实际使用方式。

Dask-ML 的三大应用方式

  1. 扩展 Scikit-Learn
    • 提供了许多与Scikit-Learn API兼容的并行算法。
  2. 增量学习 (核外学习)
    • 支持在无法一次性加载的数据上逐步训练模型。
  3. 可扩展的预处理
    • 提供并行 StandardScaler 等工具,可由训练分区样本逐特征归约均值与标准差。

实践:多分区 Dask-ML 线性回归接口演示

Dask-ML实践: 目标

回到之前检查的确定 large_transactions.csv 数据集,尝试训练一个线性回归模型,预测交易量 volume

我们将用 dask-ml 演示多分区接口;该约 40MB 合成文件通常可一次性载入内存,因此不能作为核外或超内存证据。本章的真实数据验证见后文恒瑞医药分钟数据版本。

步骤1 - 准备数据

首先,我们像之前一样用Dask读取数据,并定义我们的特征(X)和目标(y)。

代码
# 为“步骤1 - 准备数据”导入 `dask.dataframe` 并绑定 `dd`,用于建立当前任务的分区表格与惰性任务图。
import dask.dataframe as dd
# 为“步骤1 - 准备数据”,导入 Dask-ML 缩放器,逐特征归约全量机制拟合样本的均值与标准差;本页不作样本外验证。
from dask_ml.preprocessing import StandardScaler
# 为“步骤1 - 准备数据”,导入 Dask-ML 线性回归器;本章显式确定 L2 惩罚、C=1.0 与 ADMM 求解器,不把它称为普通 OLS。
from dask_ml.linear_model import LinearRegression

# 读取数据
ddf = dd.read_csv('large_transactions.csv', blocksize='8MB')

# 为“步骤1 - 准备数据”,定义特征和目标
# 为了简化,我们只用 price 和 stock_id 作为特征
# 注意:直接使用 stock_id 是不好的做法,这里仅为演示
X = ddf[['price', 'stock_id']]
# 提取示例目标数组,供 Dask-ML 缩放与线性回归。
y = ddf['volume']

# 把两列分区特征转换为保留分块信息的 Dask Array,供 Dask-ML 估计器消费。
X = X.to_dask_array(lengths=True)
# 把目标 Series 转换为与特征行分区对齐的 Dask Array,而非物化为 NumPy 数组。
y = y.to_dask_array(lengths=True)

# 报告特征与目标仍保持为分块的 Dask Array,而非已物化数组。
print("X 和 y 已准备为 Dask Arrays.")
# 输出 Dask Array 的分块与形状表示,确认特征没有被物化成 NumPy 数组。
print(X)
X 和 y 已准备为 Dask Arrays.
dask.array<read-_to_string_dtype-getitem-values, shape=(1000000, 2), dtype=float64, chunksize=(202096, 2), chunktype=numpy.ndarray>

步骤2 - 数据标准化

  • fit:跨分区归约拟合样本的均值与方差,并物化拟合属性。
  • transform:使用已拟合统计量,返回的 Dask Array 仍保持惰性。
  • 检查点:统计量只能来自机制拟合样本。
代码
# 创建将由机制拟合样本逐特征估计统计量的 Dask 缩放器。
scaler = StandardScaler()

# 逐特征汇总机制拟合样本的分区统计量,并物化均值与尺度属性。
scaler.fit(X)
# 使用确定的拟合样本逐特征统计量构建惰性标准化数组。
X_scaled = scaler.transform(X)

# 打印已物化均值属性与惰性转换结果的类型,核对 `fit` 和 `transform` 的执行边界。
print({'均值属性类型': type(scaler.mean_).__name__, '转换结果类型': type(X_scaled).__name__})  # 区分已物化拟合属性与惰性转换集合
# 打印惰性 Dask Array 的表示、形状与分块元数据;此句不计算列均值或验证缩放数值。
print(X_scaled)
{'均值属性类型': 'ndarray', '转换结果类型': 'Array'}
dask.array<truediv, shape=(1000000, 2), dtype=float64, chunksize=(202096, 2), chunktype=numpy.ndarray>

步骤3 - 训练模型

  • 固定规格penalty='l2'C=1.0solver='admm'
  • 模型含义:在 Dask-ML 2025.1.0 中,这是 \(\lambda=1/C\) 的 L2 惩罚线性回归,不是普通最小二乘。
  • 执行边界:是否跨分区仍取决于输入、solver、scheduler 与实际使用。
代码
# 显式确定 L2 惩罚、C=1.0 与 ADMM 求解器,避免类默认值随版本变化并误称为普通 OLS。
lr = LinearRegression(penalty='l2', C=1.0, solver='admm')

# 训练模型
# 这一步会触发计算!
# Dask-ML会在后台读取数据块,进行计算
print('开始训练模型...')
# `.fit()` 触发分区计算并物化线性回归系数与截距,完成本页机制样本拟合。
lr.fit(X_scaled, y)
# 在拟合完成后报告模型参数已由分区归约计算得到。
print('模型训练完毕!')

# 输出本例确定的模型族、正则强度与求解器,使结果不依赖未报告默认值。
print({'模型族': 'L2惩罚线性回归', 'penalty': lr.penalty, 'C': lr.C, 'solver': lr.solver})
# 打印模型系数
print(f'模型系数 (Coefficients): {lr.coef_}')
# 打印已拟合线性回归的截距,核对 `.fit()` 已物化模型参数而非返回惰性预测集合。
print(f'模型截距 (Intercept): {lr.intercept_}')
开始训练模型...
模型训练完毕!
{'模型族': 'L2惩罚线性回归', 'penalty': 'l2', 'C': 1.0, 'solver': 'admm'}
模型系数 (Coefficients): [ 4.93338112 -2.6660765 ]
模型截距 (Intercept): 5044.245107269162

解读 Dask-ML 的 .fit()

  • 当我们调用.fit()时,Dask-ML和Dask在后台协同工作。
  • Dask负责高效地从磁盘读取数据块(partitions)。
  • Dask-ML则负责在这些数据块上执行并行的机器学习算法(例如,计算每个块的梯度)。
  • 最后,Dask-ML负责聚合所有块的结果,得到一个统一的最终模型。
  • 整个过程对用户来说是透明的,API与Scikit-Learn高度一致。

结论:关键思想的转变

阶段小结 (1): 思维的转变

从“一次性加载”到“分块处理”: 这是从传统数据分析迈向大数据分析最根本的思维转变。

你必须放弃“所有数据都在我眼前”的假设,开始用“流水线”和“分治”的眼光看待问题。

阶段小结 (2): 工具的力量

懒惰执行的力量: 通过先规划(构建任务图)再按需物化(如 .compute().persist().head().fit()),Dask 可优化执行并在合适分区、算法、调度和 spill 配置下处理超出单机内存的数据;本章小文件不构成该能力的实测证据。

这是一种更智能、更具伸缩性的计算范式。

阶段小结 (3): 生态的威力

Python生态的威力: Dask和PySpark让Python这门通用语言具备了处理工业级大数据的能力,而Dask-ML则将这种能力无缝延伸到了机器学习领域。

你所学的技能,可以直接应用于解决真实世界的大规模问题。

课堂核心独立任务

检查:何时不应调用 .compute()

  1. 场景一: 你需要处理一个2GB的CSV文件,并且你的电脑有16GB内存。你会选择哪个库(Pandas还是Dask)?为什么?
  2. 场景二: 你需要处理一个5TB的用户行为日志数据集,并构建一个复杂的推荐系统,你的公司有一个大型计算集群。你应该考虑使用哪个库(Dask还是PySpark)?为什么?
  3. 场景三: 一个 Dask groupby 仍是 lazy collection。请分别判断何时可以调用 .compute()、何时必须拒绝:结果预计 2KB 且任务图峰值内存 3GB;结果大小未知且 shuffle 峰值可能超过 16GB。

反馈:文件大小不能单独决定框架

框架选择

  • 2GB CSV:依据解析后的峰值内存与墙钟基准决定 Pandas 是否可行。
  • 5TB + 复杂 shuffle + 大集群:优先评估 PySpark,并核对容错与团队运维能力。

反馈:何时可以物化?

  • 可以:结果约 2KB、任务图峰值约 3GB,且运行时间和内存均低于预算;同时记录墙钟与峰值。
  • 拒绝:结果大小未知或 shuffle 峰值可能超预算;先做分区抽样、分阶段归约、persist 或写出分区结果。
  • 订正路径:混淆惰性与任务图,回看惰性执行任务图;遗漏物化副作用,回看执行与物化

独立任务:Dask 分区归约

使用课程确定的 large_transactions.csv,完成以下任务:

  • 先报告文件大小、字段、行数与实际 npartitions;数据或分区不符即停止。
  • timestamp列转换为年份。
  • 计算每年最高最低交易价格(price)。

实现提示dd.to_datetime().dt.yeargroupby()agg()

独立任务:复算、资源与失败路径

  • 用 Pandas 对同一数据版本、同一字段与同一年度口径复算,并以 atol=0.001 核对 Dask 结果。
  • 分别记录 Dask 与 Pandas 的墙钟秒数和进程峰值内存;当前 41MB 数据版本只比较接口与口径,不证明超内存优势。
  • 保留失败路径:哈希/字段/行数不符、分区数不大于 1、两框架年度结果不一致时,明确停止且不得提交性能结论。

习题 4:Dask 合并与框架选择

假设你还有一个stocks_info.csv文件,包含stock_idstock_name两列。

  • 合并设计:说明如何用 Dask 按 stock_idlarge_transactions.csv 合并。
  • 大小分支:分别说明 stocks_info.csv 为小表或大表时的策略。
  • 结果用途:按 stock_name 分析交易情况。
  • 框架复答:把选择与峰值内存、墙钟、分区/shuffle、容错及团队运维逐项对应。

练习:先完成再查看答案:年度归约作业结果

  • 计算证据:CSV 哈希与分区、年度极值、Pandas 同口径基线、墙钟与峰值内存。
  • 合并证据:小表/大表策略、stock_id 唯一性、合并前后行数与未匹配比例。
  • 决策证据:两项框架场景自查和失败路径。
  • 重做条件:缺任一完成内容、缺分区,或只报速度不报口径。

课堂核心完整反馈(1/5):场景自查

  • 场景一判断:2GB 文件在 16GB 内存中未必安全;解析、字符串列与中间副本会放大内存。
  • 场景一证据:先做 Pandas 分列、分块基准;峰值不可控时再评估 Dask。
  • 场景二默认:5TB 日志、复杂推荐系统与大型集群下,优先选择 PySpark。
  • 改选条件:已有 Dask Distributed 运维能力,且分区、倾斜、shuffle 压测与容错演练均达标。

4 分评分:题面一致、分区/shuffle、容错/团队、资源与失败恢复计划各 1 分;只写框架名或只按文件大小判断不得分。

课堂核心完整反馈(2/5):用相同口径比较 Pandas 与 Dask

代码
import time

# Pandas:一次读入所需两列,再计算各年的最低价与最高价。
pandas_start = time.perf_counter()
pandas_frame = pd.read_csv(file_path, usecols=['timestamp', 'price'])
pandas_frame['year'] = pd.to_datetime(pandas_frame['timestamp']).dt.year
pandas_price_bounds = pandas_frame.groupby('year')['price'].agg(['min', 'max']).sort_index()
pandas_seconds = time.perf_counter() - pandas_start

# Dask:把同一文件分成多个部分,最后用 `.compute()` 得到结果。
dask_start = time.perf_counter()
dask_frame = dd.read_csv(file_path, usecols=['timestamp', 'price'], blocksize='8MB')
dask_frame['year'] = dd.to_datetime(dask_frame['timestamp']).dt.year
annual_price_bounds = dask_frame.groupby('year')['price'].agg(['min', 'max']).compute().sort_index()
dask_seconds = time.perf_counter() - dask_start

assert dask_frame.npartitions > 1
runtime_comparison = pd.DataFrame({
    '工具': ['Pandas', 'Dask'],
    '运行时间(秒)': [pandas_seconds, dask_seconds]
})
display(pd.DataFrame([{'文件大小(字节)': os.path.getsize(file_path), 'Dask 分区数': dask_frame.npartitions}]))
display(runtime_comparison.round(3))
表 1: 同一份 CSV 的分区数与两种工具运行时间
文件大小(字节) Dask 分区数
0 42406725 5
工具 运行时间(秒)
0 Pandas 3.120
1 Dask 8.591

课堂核心完整反馈(3/5):年度极值一致性

  1. Dask 参考代码
代码
pd.testing.assert_frame_equal(annual_price_bounds, pandas_price_bounds, check_exact=False, atol=0.001)  # 两框架必须同口径一致
display(annual_price_bounds)  # 输出可检查的完整答案
print({'与 Pandas 一致': True, '允许的数值误差': 0.001})
min max
year
2017 10.0 500.0
2018 10.0 500.0
2019 10.0 500.0
2020 10.0 500.0
{'与 Pandas 一致': True, '允许的数值误差': 0.001}

确认输入包含五列字段和一百万行后,以上动态四行表就是期望输出;价格按两位小数核对,绝对容差为 0.001

课堂核心完整反馈(4/5):归约解释与失败路径

  • 归约机制:每个年份先接收各源分区的局部极值,再做跨分区归约。
  • 执行范围.compute() 触发完整年度分组任务图;最终四行不代表只读取少数分区。
  • 失败路径:文件、字段或行数不符先停止;年度结果差异超过 0.001 时回查日期与分组步骤。
  • 证据边界:当前运行时间只适用于这台电脑和教学数据,不证明中国市场、超内存或集群性能。

课堂核心完整反馈(5/5):合并验证

  1. 合并答案:若 stocks_info.csv 很小,可先读成单分区并广播式合并;若同样很大,应按 stock_id 合理分区并检查 shuffle 成本。合并后验证行数、主键唯一性和未匹配比例,防止多对多连接放大样本。

公开中国数据练习:Dask 分区计算

代码
from pathlib import Path  # 管理公开 A 股分钟数据路径
from urllib.request import urlretrieve  # 复用本章已安装的浏览器标识下载器
import dask.array as da  # 建立惰性分区数组并执行聚合
import h5py  # 访问非 Pandas HDF5 字段视图
# 按 Linux 共享数据、Windows 共享数据、项目缓存的顺序选择分钟行情文件。
minute_path = next((candidate_path for candidate_path in [Path('/home/ubuntu/r2_data_mount/data/stock/minbar/600276.XSHG.h5'), Path('C:/qiufei/data/stock/minbar/600276.XSHG.h5'), Path('data/course/600276.XSHG.h5')] if candidate_path.exists()), Path('data/course/600276.XSHG.h5'))
if not minute_path.exists():
    minute_path.parent.mkdir(parents=True, exist_ok=True)
    urlretrieve('https://assets.qiufei.site/data/stock/minbar/600276.XSHG.h5', minute_path)
minute_store = h5py.File(minute_path, 'r')  # 保持文件打开以支持 Dask 延迟读取
close_field = minute_store['data'].fields('close')  # 只暴露收盘价字段而不加载整个结构化数组
minute_record_count = int(close_field.shape[0])  # 在关闭文件前保存总记录数
partitioned_close = da.from_array(close_field, chunks=200000)  # 按二十万行构造多个惰性分区
partition_count = partitioned_close.npartitions  # 记录任务图中的实际分区数
assert partition_count > 1, '分区数必须大于1才能验证分区计算'  # 防止单分区伪装为规模化演示
# 无论均值计算成功与否,都在退出本页计算时关闭 HDF5 句柄。
try:
    mean_close = float(partitioned_close.mean().compute())  # 触发跨分区均值计算并返回标量
# 用 finally 保证计算异常不会遗留打开的公开数据文件。
finally:
    minute_store.close()  # 释放对分钟数据版本的文件句柄
print({'文件': minute_path.name, '数据集': 'data.close', '记录数': minute_record_count, '分区数': partition_count, '平均收盘价': round(mean_close, 4)})  # 输出实际分区与计算结果
表 2: 恒瑞医药分钟 HDF5 的 Dask 分区与实际均值计算
{'文件': '600276.XSHG.h5', '数据集': 'data.close', '记录数': 1227600, '分区数': 7, '平均收盘价': 46.3662}

公开真实数据练习:Dask-ML 下一分钟收益基线

  • 输入规模:恒瑞医药分钟数据的前 300,001 条记录。
  • 特征与目标:当前分钟收盘价、振幅、对数成交量;目标为下一分钟简单收益。
  • 时间切分:前 240,000 个样本训练,后 60,000 个样本最终测试。
  • 证据边界:只验证 Dask-ML 接口与时间切分,不声称超内存能力,也不作交易建议。
代码
import time  # 记录真实多分区模型的墙钟时间
from dask_ml.preprocessing import StandardScaler as DaskStandardScaler  # 在训练分区逐特征估计均值与标准差
from dask_ml.linear_model import LinearRegression as DaskLinearRegression  # 拟合显式确定规格的真实分钟收益 L2 线性模型
from sklearn.metrics import mean_squared_error  # 统一模型与零收益基线的最终测试指标
local_model_store = h5py.File(minute_path, 'r')  # 保持 HDF5 打开以支持惰性字段读取
local_model_count = 300001  # 确定教学计算规模并保留 300000 个监督样本
local_close = da.from_array(local_model_store['data'].fields('close'), chunks=50000)[:local_model_count]  # 惰性读取收盘价字段
local_high = da.from_array(local_model_store['data'].fields('high'), chunks=50000)[:local_model_count]  # 惰性读取最高价字段
local_low = da.from_array(local_model_store['data'].fields('low'), chunks=50000)[:local_model_count]  # 惰性读取最低价字段
local_volume = da.from_array(local_model_store['data'].fields('volume'), chunks=50000)[:local_model_count]  # 惰性读取成交量字段
local_features = da.stack([local_close[:-1], (local_high[:-1] - local_low[:-1]) / local_close[:-1], da.log1p(local_volume[:-1])], axis=1)  # 组合当分钟收盘价、振幅与对数成交量,预测下一分钟收益。
local_target = local_close[1:] / local_close[:-1] - 1  # 构造严格向前一期的简单收益目标
local_training_end = 240000  # 确定前 80% 训练与后 20% 最终测试边界
local_training_features = local_features[:local_training_end].rechunk((40000, 3))  # 对齐六个训练分区及完整特征轴
local_training_target = local_target[:local_training_end].rechunk(40000)  # 与训练特征使用相同分区边界
local_testing_features = local_features[local_training_end:].rechunk((30000, 3))  # 对齐两个最终测试分区
local_testing_target = local_target[local_training_end:].rechunk(30000)  # 与最终测试特征使用相同分区边界
代码
local_start_time = time.perf_counter()  # 从实际拟合前开始计时
local_dask_scaler = DaskStandardScaler().fit(local_training_features)  # 只从训练分区样本逐特征估计均值与标准差
local_training_scaled = local_dask_scaler.transform(local_training_features)  # 用训练样本逐特征统计量转换训练分区
local_testing_scaled = local_dask_scaler.transform(local_testing_features)  # 用确定的训练样本逐特征统计量转换最终测试分区
local_dask_model = DaskLinearRegression(penalty='l2', C=1.0, solver='admm').fit(local_training_scaled, local_training_target)  # 在真实训练分区拟合与前例一致的 L2 惩罚线性模型
local_model_predictions = local_dask_model.predict(local_testing_scaled).compute()  # 触发最终测试分区预测计算
local_target_values = local_testing_target.compute()  # 只在最终计分时收集最终测试目标
local_model_seconds = time.perf_counter() - local_start_time  # 记录拟合与最终测试预测墙钟时间
local_zero_predictions = np.zeros_like(local_target_values)  # 构造同窗零收益基线
local_dask_results = pd.DataFrame({'模型': ['Dask-ML L2线性回归', '零收益基线'], '规格': ['penalty=l2; C=1.0; solver=admm', '固定预测=0'], '最终测试MSE': [mean_squared_error(local_target_values, local_model_predictions), mean_squared_error(local_target_values, local_zero_predictions)]})  # 同窗报告确定模型规格与最终测试误差,避免把惩罚回归误称为 OLS
local_model_store.close()  # 完成惰性计算后关闭真实 HDF5 文件
print({'训练样本': local_training_end, '最终测试样本': len(local_target_values), '训练分区': local_training_features.npartitions, '最终测试分区': local_testing_features.npartitions, '运行秒': round(local_model_seconds, 2)})  # 输出任务规模与资源边界
display(local_dask_results.round(10))  # 如实展示模型是否胜过零收益基线
{'训练样本': 240000, '最终测试样本': 60000, '训练分区': 6, '最终测试分区': 2, '运行秒': 45.34}
表 3: 恒瑞医药真实分钟数据的 Dask-ML 与零收益基线最终测试比较
模型 规格 最终测试MSE
0 Dask-ML L2线性回归 penalty=l2; C=1.0; solver=admm 0.000004
1 零收益基线 固定预测=0 0.000004

结果不理想时如何解释:若 Dask-ML 未胜零收益基线,只能写“当前固定分钟样本与线性规格没有增量预测价值”。该实验的贡献是验证真实 HDF5 的多分区读取、惰性预处理、拟合与最终测试计分链,而不是证明超内存能力或市场可预测性。

本章小结

  • 能把选择性读取、分区变换、惰性任务图、归约与一次最终测试评分串成可检查链。
  • 是否 .compute() 取决于结果大小、内存预算与副作用;分区数需由数据规模和调度开销验证。
  • 当前实验不证明超内存实际使用能力或市场可预测性。
  • 下一章综合资产定价与机器学习,并继续执行点时、基线和失败报告说明。