4  数据读取与存储:从本地到大数据系统

4.1 引言与学习目标

数据分析的第一步是将数据从各种来源读入 Python 环境。在真实的金融数据分析工作中,数据可能存储在 CSV 文件、HDF5 文件、Parquet 文件、SQL 数据库,甚至是网络 API 中。选择合适的数据格式和读写方法,不仅影响代码的运行效率,还直接决定了能否在有限的内存中处理大规模数据集。

学习目标

完成本章后,你应能:

  • 根据按行读取、列投影、谓词过滤、跨工具交换和人工检查需求,比较 CSV、HDF5、Parquet 与 Arrow 的适用边界;
  • 识别 DataFrame 的行、列与索引,完成列选择、布尔行筛选、索引设置以及 groupby().agg() 最小聚合;
  • 使用 read_csv()dtypeusecolschunksize,并用全量安全子集核对分块聚合结果;
  • 创建 SQLite 表、写出聚合 SQL,并把查询结果与等价 Pandas 聚合进行一致性检查;
  • 检查 HDF5 节点是 table 还是 fixed,确认查询列属于 data_columns 后再使用 where
  • 解释 Pandas 与 Arrow 转换、Arrow IPC 序列化和跨工具零拷贝分别在哪些条件下可能发生数据复制。

目标—活动—核心练习/答案映射

正式目标 正文活动 核心评价证据
比较 CSV、HDF5、Parquet 与 Arrow 小节 4.2小节 4.5小节 4.6 的格式决策 练习 4.2、4.4 及答案中的同源写出比较和条件判断
使用最小 DataFrame 操作 小节 4.3 的列、行、索引与分组桥接 练习 4.3、4.5 及答案中的选择、聚合、排序与逐项核对
正确使用 dtypeusecolschunksize 小节 4.2.1小节 4.3.1 的参数化读取 练习 4.5 及答案中的 schema 断言、分块聚合和全量安全子集核验
对照 SQLite 与 Pandas 聚合 小节 4.4.1小节 4.4.2 的等价查询 练习 4.3 及答案中的逐键、逐指标一致性断言
检查 HDF5 布局与可查询列 小节 4.5.2 的节点元数据检查 练习 4.1、4.4 及答案中的安全行筛选、列投影和条件说明
解释 Arrow 转换和复制条件 小节 4.5.4 的缓冲区与 IPC 活动 练习 4.4 及答案中的三个反绝对化判断

前置知识

读者应会使用第 2 章的循环与文件路径,以及第 3 章的 dtype 和数组内存概念;不要求预先掌握 Pandas DataFrame。本章在首次核心练习前先桥接 DataFrame 的行、列、索引、选择和分组聚合,再讲 CSV 与数据库 I/O。SQL 小节从 SELECTGROUP BY 和聚合函数的最小例开始,不要求数据库先修。

4.2 结构化文本读写(CSV)的进阶参数与分块迭代

4.2.1 CSV 读取的核心参数

pd.read_csv() 是 Pandas 中最常用的数据读取函数。除了基本的文件路径参数外,掌握以下进阶参数对于处理真实世界的数据至关重要:

import pandas as pd  # 导入 Pandas 用于数据操作
import platform  # 导入平台检测模块
from pathlib import Path  # 使用路径对象管理本章全部临时教学文件
from tempfile import TemporaryDirectory  # 创建内核结束时自动清理的临时目录
DATA_ROOT = 'C:/qiufei/data' if platform.system() == 'Windows' else '/home/ubuntu/r2_data_mount/data'  # 根据操作系统选择数据根路径
chapter_io_tmpdir = TemporaryDirectory()  # 保持对象存活,使后续代码块共享同一临时目录
chapter_io_root = Path(chapter_io_tmpdir.name)  # 建立跨代码块复用的临时路径根
demo_csv_path = chapter_io_root / 'demo_trades.csv'  # 指定CSV教学副本的临时路径
demo_sqlite_path = chapter_io_root / 'demo_stock.db'  # 指定SQLite教学数据库的临时路径
demo_parquet_path = chapter_io_root / 'demo_trades.parquet'  # 指定Parquet教学副本的临时路径

教师提供的数据准备(只运行核对,不评分)

下一块用尚未正式桥接的 HDF5 读取与 DataFrame 整理 API,把海康威视真实行情转换为本章统一的 trade_data 和临时 CSV。学生只需运行并核对公司代码、列名与记录数;该块不要求解释、修改或默写,也不计入核心评分。

stock_price_path = f'{DATA_ROOT}/stock/stock_price_pre_adjusted.h5'  # 指向本地前复权日行情表
hikvision_price_df = pd.read_hdf(  # 从可查询 HDF5 table 中选择海康威视
    stock_price_path,  # 复用统一数据路径,避免硬编码平台差异
    where="order_book_id == '002415.XSHE'"  # 通过已索引的证券代码列下推行过滤
).reset_index()  # 将底表索引恢复为列以便导出通用 CSV
required_columns = {'order_book_id', 'date', 'close', 'volume'}  # 声明教学副本所需的数据合同
assert required_columns.issubset(hikvision_price_df.columns)  # 写出前验证关键列全部存在

trade_data = hikvision_price_df.loc[  # 只保留后续 CSV 与 SQL 示例需要的四列
    :, ['date', 'order_book_id', 'close', 'volume']  # 投影日期、代码、价格与成交量
].rename(columns={  # 把本地 schema 映射为本章统一教学 schema
    'date': 'trade_time',  # 将底表日期列统一命名为交易时间
    'order_book_id': 'stock_code',  # 将证券标识统一命名为股票代码
    'close': 'price'  # 将收盘价统一命名为价格
})  # 完成教学副本字段标准化
trade_data.to_csv(demo_csv_path, index=False)  # 把CSV教学副本写入临时目录而不改写仓库data文件
print(f'已导出 {len(trade_data):,} 行海康威视真实行情')  # 报告教学副本的真实记录数
trade_data.head()  # 预览前5行
已导出 3,789 行海康威视真实行情
trade_time stock_code price volume
0 2010-05-28 002415.XSHE 3.4389 613770372.0
1 2010-05-31 002415.XSHE 3.5678 310101714.0
2 2010-06-01 002415.XSHE 3.6622 197904240.0
3 2010-06-02 002415.XSHE 3.6139 170932752.0
4 2010-06-03 002415.XSHE 3.4251 136089396.0

4.3 DataFrame 最小桥接:行、列、索引与聚合

DataFrame 是带列名的二维表。每一列保存一种变量,每一行保存一条观测,index 则是定位行的标签。后续 I/O 函数返回的主要就是这种对象,因此先掌握四个最小动作:查看结构、选择列、按条件选择行、设置索引。方括号内放一个列名会返回一维 Series;放列名列表会保持二维 DataFrame.loc[行条件, 列列表] 同时表达行筛选与列投影。

assert trade_data.columns.tolist() == ['trade_time', 'stock_code', 'price', 'volume']  # 核对教师交付表的四列合同
price_series = trade_data['price']  # 用单列名取得一维收盘价Series
price_volume_frame = trade_data[['price', 'volume']]  # 用列名列表保持二维DataFrame
valid_trade_rows = trade_data.loc[  # 同时按布尔条件选行并按列名投影
    trade_data['price'].notna() & trade_data['volume'].gt(0),  # 只保留价格非缺失且有成交量的真实观测
    ['trade_time', 'stock_code', 'price', 'volume'],  # 明确输出列及其顺序
]  # 完成行列选择
trades_by_time = valid_trade_rows.set_index('trade_time').sort_index()  # 把交易时间设为有序行索引
assert price_series.ndim == 1 and price_volume_frame.ndim == 2  # 核对单列与列列表返回对象维数不同
trades_by_time.head()  # 展示以真实交易日为索引的前五条观测
stock_code price volume
trade_time
2010-05-28 002415.XSHE 3.4389 613770372.0
2010-05-31 002415.XSHE 3.5678 310101714.0
2010-06-01 002415.XSHE 3.6622 197904240.0
2010-06-02 002415.XSHE 3.6139 170932752.0
2010-06-03 002415.XSHE 3.4251 136089396.0

分组聚合先用 groupby('分组键') 把同行业或同证券记录归组,再用 agg(输出列名=('输入列', '函数')) 为每组计算指标。reset_index() 把分组键从索引恢复为普通列,便于写出文件或与 SQL 查询逐列比较。

security_summary_bridge = (  # 对教师提供的真实单证券表演示命名聚合语法
    valid_trade_rows.groupby('stock_code', observed=True)  # 按证券代码形成互斥记录组
    .agg(record_count=('price', 'size'), average_price=('price', 'mean'), total_volume=('volume', 'sum'))  # 为每组命名三项输出指标
    .reset_index()  # 把证券代码恢复为普通列,形成可写出的二维表
)  # 完成最小分组聚合
assert security_summary_bridge['record_count'].sum() == len(valid_trade_rows)  # 核对分组后记录数守恒
security_summary_bridge  # 展示后续SQL对照会复用的结果形状
stock_code record_count average_price total_volume
0 002415.XSHE 3735 21.658679 1.468341e+11

关键参数详解

parsed_trades = pd.read_csv(  # 使用进阶参数读取CSV
    demo_csv_path,  # 读取上一代码块写入同一临时目录的真实行情教学副本
    parse_dates=['trade_time'],  # 自动解析日期时间列
    dtype={'stock_code': 'category'},  # 将股票代码读取为分类类型,节省内存
    usecols=['trade_time', 'stock_code', 'price']  # 只读取需要的列,跳过volume列
)  # 完成带显式日期、类型和列投影参数的 CSV 读取
print(f'数据类型:\n{parsed_trades.dtypes}')  # 验证各列的数据类型
print(f'\n内存占用: {parsed_trades.memory_usage(deep=True).sum() / 1024:.1f} KB')  # 查看内存占用
数据类型:
trade_time    datetime64[ns]
stock_code          category
price                float64
dtype: object

内存占用: 63.1 KB
参数 作用 典型使用场景
parse_dates 自动解析日期列 时间序列数据
dtype 指定列的数据类型 内存优化、加速读取
usecols 只读取指定列 宽表中只需部分字段
nrows 只读取前N行 快速预览大文件
chunksize 分块迭代读取 内存无法容纳的大文件
encoding 指定文件编码 中文数据常用 utf-8gbk
na_values 自定义缺失值标记 数据源使用 --N/A 表示缺失

4.3.1 大文件的分块迭代读取

当 CSV 文件超过内存容量时,chunksize 参数是最实用的解决方案。它将文件分成若干小块(Chunk),逐块读入内存进行处理:

total_trade_volume = 0  # 初始化成交量累计值为0
chunk_count = 0  # 初始化已处理的块数为0
chunk_reader = pd.read_csv(demo_csv_path, chunksize=2000)  # 从同一临时CSV路径每次读取2000行
for chunk in chunk_reader:  # 逐块遍历文件
    total_trade_volume += chunk['volume'].sum()  # 累加当前块的成交量
    chunk_count += 1  # 块计数加1
print(f'共处理 {chunk_count} 个数据块')  # 输出总共处理了多少块
print(f'总成交量: {total_trade_volume:,}')  # 输出成交量总和(带千位分隔符)
共处理 2 个数据块
总成交量: 146,834,111,963.75

选择合适的 chunksize

chunksize 的选择需要在内存占用和处理效率之间权衡:

  • 太小(如 100 行):过多的循环迭代次数导致 Python 解释器开销增大
  • 太大(如 5000 万行):可能超出可用内存
  • 推荐:根据单行的内存占用和可用内存计算,通常 1 万到 100 万行是合理范围

4.4 关系型数据库的连接与 SQL 查询

在企业级数据分析中,数据通常存储在关系型数据库中(如 MySQL、PostgreSQL、SQLite)。Pandas 通过 SQLAlchemy 库提供了与数据库的无缝集成。

4.4.1 SQLite 本地数据库示例

SQLite 是一个轻量级的嵌入式数据库,非常适合学习和原型开发:

import sqlite3  # 导入Python内置的SQLite数据库驱动
connection = sqlite3.connect(demo_sqlite_path)  # 在本章临时目录创建SQLite教学数据库
trade_data_for_db = pd.read_csv(demo_csv_path)  # 从同一临时CSV路径读取真实教学副本
trade_data_for_db.to_sql('trades', connection, if_exists='replace', index=False)  # 将数据写入数据库的trades表
print('数据已成功写入 SQLite 数据库')  # 确认写入完成
数据已成功写入 SQLite 数据库
# 把分组统计查询保存为可复用文本,确保 SQL 逻辑与执行入口分离
sql_query = """
    SELECT stock_code,
           COUNT(*) as trade_count,
           AVG(price) as avg_price,
           SUM(volume) as total_volume
    FROM trades
    GROUP BY stock_code
    ORDER BY total_volume DESC
"""  # 编写SQL查询:按股票代码分组,统计交易次数、均价和总成交量
stock_summary_from_sql = pd.read_sql(sql_query, connection)  # 执行SQL查询并返回DataFrame
stock_summary_from_sql  # 展示按成交量降序排列的统计结果
stock_code trade_count avg_price total_volume
0 002415.XSHE 3789 21.505153 1.468341e+11
connection.close()  # 关闭数据库连接,释放资源
print('数据库连接已关闭')  # 确认连接关闭
数据库连接已关闭

4.4.2 SQL 查询 vs Pandas 操作

SQL 查询和 Pandas 操作在功能上高度对应:

SQL 操作 Pandas 等价操作
SELECT col1, col2 df[['col1', 'col2']]
WHERE condition df[df['col'] > value]
GROUP BY col df.groupby('col')
ORDER BY col DESC df.sort_values('col', ascending=False)
JOIN ... ON pd.merge(df1, df2, on='key')
COUNT(*) df.groupby('col').size()

在后续的 小节 13.1 中,我们将学习 DuckDB——一种可以在进程内直接对 Pandas DataFrame 和文件执行 SQL 查询的高性能分析引擎。

4.5 文本记录、列式文件与列式内存

4.5.1 CSV 文本记录与 Parquet 列式文件

选择数据格式前,必须先区分文本序列化、磁盘文件布局与内存布局。CSV、Parquet 和 Arrow 分别处在不同层次,不能简单并称为三种可互换的“存储方式”。

CSV:分隔符文本记录

CSV 按记录写出字段的文本表示,换行通常分隔记录,分隔符区分字段。它没有统一的物理类型、页索引或列块元数据;即使 usecols 只返回少数列,解析器通常仍要扫描并解析相关文本字节。因此,CSV 便于人工检查和跨系统交换,却不能仅凭“按行书写”推断随机行读取一定高效。

Parquet:带层级元数据的列式文件格式

Parquet 文件由一个或多个 row group 组成;每个 row group 内,每列对应 column chunk,column chunk 再由 page 组成。读取引擎可跳过未投影列的 column chunk,并在统计信息和谓词支持时跳过不相关 row group。压缩效果和过滤收益取决于类型、编码、列基数、row group 布局与查询,不能宣称所有数据都必然更小或更快。

Arrow 则定义列式内存数组、缓冲区和 schema;持久化或跨进程传输需使用 Arrow IPC、Feather 等具体协议或文件格式。后文 小节 4.5.4 单独讨论其复制与序列化边界。

图 4.1 展示本节讨论对象的可视化结果,读图时应结合正文给出的口径与限制。

CSV 文本记录 vs Parquet 列式文件 文本记录 (CSV) 代码 价格 成交量 日期 600276 45 5200 01-02 002415 38.5 8100 01-02 600104 15 3400 01-02 字段以文本形式写入记录 列投影通常仍需扫描文本 列式存储 (Parquet) 代码 600276 002415 600104 价格 45 38.5 15 成交量 5200 8100 3400 row group 内按列形成 chunk 投影与过滤收益取决于布局 性能必须在同一数据与环境下测量 指标 CSV Parquet 文件大小 取决于文本表示 取决于编码与压缩 读取全部 需同机实测 需同机实测 读取3列 解析时仍扫描文本 列投影可减少 I/O
图 4.1: CSV 文本记录与 Parquet 列式文件的逻辑组织差异;图中数值仅为假设记录,性能没有硬编码倍数

4.5.2 HDF5 格式深度解析

HDF5(Hierarchical Data Format 5)是一种支持分层数据组织的文件格式,特别适合存储大型科学数据集。本教材的核心数据集均采用 HDF5 格式。

HDF5 的核心优势

  1. 分层结构:类似文件系统,一个 HDF5 文件内部可以包含多个”组”和”数据集”
  2. 选择性读取有前提:Pandas 只有在节点以 format='table' 写入,且筛选列被声明为 data_columns 时,才能用 where 下推相应条件;fixed 节点不支持这种查询
  3. 压缩支持:支持多种压缩算法(如 blosc、zlib),大幅减小文件体积
  4. 类型丰富:原生支持多维数组、字符串、日期时间等类型
import pandas as pd  # 导入 Pandas 用于 HDF5 读取
stock_price_path = f'{DATA_ROOT}/stock/stock_price_pre_adjusted.h5'  # 构建前复权股价数据文件路径
import pandas as pd  # 导入 Pandas
stock_price_path = f'{DATA_ROOT}/stock/stock_price_pre_adjusted.h5'  # 构建文件路径
store = pd.HDFStore(stock_price_path, mode='r')  # 以只读模式打开HDF5文件
print('HDF5 文件中的数据集:')  # 提示输出内容
print(store.keys())  # 列出文件中所有数据集的键名
stock_storer = store.get_storer('/data')  # 取得行情节点元数据以检查写入布局
print(f'节点布局: {stock_storer.format_type}')  # 验证本地行情节点采用可查询的 table 布局
print(f'可查询列: {stock_storer.data_columns}')  # 列出允许 where 下推的 data_columns
store.close()  # 关闭文件句柄释放资源
HDF5 文件中的数据集:
['/data']
节点布局: table
可查询列: ['date', 'order_book_id']

选择性读取:只加载需要的数据子集

对本地大行情表,优先先检查节点布局与查询列,再按研究对象选择性读取。若节点是 fixed、目标列不是 data_columns,或分析确实需要全表统计,where 就不能直接解决问题;此时应考虑重写为可查询 table、分块处理,或使用适合该算子的其他格式。

import pandas as pd  # 导入 Pandas
stock_price_path = f'{DATA_ROOT}/stock/stock_price_pre_adjusted.h5'  # 构建文件路径
# 本地节点经上一代码块确认是 table,且 order_book_id 属于 data_columns
hengrui_price_data = pd.read_hdf(stock_price_path, where='order_book_id == "600276.XSHG"')  # 仅下推读取江苏恒瑞医药真实行情
print(f'恒瑞医药股价数据维度: {hengrui_price_data.shape}')  # 输出选择性读取后的真实记录维度
hengrui_price_data.tail()  # 展示末五行以核对公司子集与字段
恒瑞医药股价数据维度: (5101, 6)
open high low close volume total_turnover
order_book_id date
600276.XSHG 2025-12-25 61.22 61.66 61.03 61.40 15518175.0 9.521223e+08
2025-12-26 61.30 61.66 60.60 61.06 21430781.0 1.311044e+09
2025-12-29 61.06 61.27 60.41 60.51 28575162.0 1.737056e+09
2025-12-30 60.40 60.60 59.90 60.18 25330461.0 1.525944e+09
2025-12-31 60.16 60.36 59.50 59.57 26827896.0 1.603949e+09

4.5.3 Parquet 格式实战

Apache Parquet 是跨分析引擎常用的列式文件格式。文件内部的 row group、column chunk 和 page 层级使列投影和基于统计信息的行组跳过成为可能;是否实际下推以及收益大小取决于写入布局、元数据和读取引擎。

trade_data_for_parquet = pd.read_csv(demo_csv_path, dtype={'stock_code': str})  # 从临时CSV读取同源交易数据并固定代码类型
trade_data_for_parquet.to_parquet(demo_parquet_path, engine='pyarrow')  # 使用PyArrow把教学Parquet写入临时目录
print('已保存为 Parquet 格式')  # 确认保存完成
已保存为 Parquet 格式
parquet_partial = pd.read_parquet(  # 从 Parquet 文件读取数据
    demo_parquet_path,  # 读取上一代码块写入同一临时目录的Parquet教学副本
    columns=['stock_code', 'price'],  # 投影下推:只读取股票代码和价格两列
    engine='pyarrow'  # 使用 PyArrow 引擎
)  # 完成 Parquet 两列投影读取
print(f'只读取了 {len(parquet_partial.columns)} 列,共 {len(parquet_partial)} 行')  # 输出读取的数据维度
parquet_partial.head()  # 预览前5行
只读取了 2 列,共 3789 行
stock_code price
0 002415.XSHE 3.4389
1 002415.XSHE 3.5678
2 002415.XSHE 3.6622
3 002415.XSHE 3.6139
4 002415.XSHE 3.4251

文件大小对比

import os  # 导入os模块用于获取文件大小
csv_file_size_kb = os.path.getsize(demo_csv_path) / 1024  # 获取临时CSV教学副本大小(KB)
parquet_file_size_kb = os.path.getsize(demo_parquet_path) / 1024  # 获取临时Parquet教学副本大小(KB)
print(f'CSV 文件大小:     {csv_file_size_kb:.1f} KB')  # 输出CSV大小
print(f'Parquet 文件大小: {parquet_file_size_kb:.1f} KB')  # 输出Parquet大小
compression_ratio = csv_file_size_kb / parquet_file_size_kb  # 计算压缩比
print(f'Parquet 压缩比:   {compression_ratio:.1f}x')  # 输出压缩比
CSV 文件大小:     153.7 KB
Parquet 文件大小: 75.8 KB
Parquet 压缩比:   2.0x

4.5.4 Apache Arrow:跨工具的内存数据标准

Apache Arrow 定义标准化的列式内存数组、缓冲区和 schema,本身不等于 Parquet 那样的磁盘文件格式。它可减少 Pandas、Polars、DuckDB、Spark 等工具之间的转换成本,但“采用 Arrow”不等于“必然零拷贝”:源与目标 dtype、缓冲区布局、分块、可空表示和所有权必须兼容。Arrow IPC 是传输或持久化 Arrow 消息的具体格式,仍需读写消息元数据与缓冲区,不能称为“无序列化开销”。

import pyarrow as pa  # 导入 PyArrow 库
import pandas as pd  # 导入 Pandas 库
sample_stock_df = pd.DataFrame({  # 创建一个简单的示例 DataFrame
    'stock_code': ['600276', '002415', '600104'],  # 使用三家长三角公司代码作为字符串标签
    'close_price': [45.0, 32.0, 15.0],  # 设置仅用于类型转换演示的假设价格
    'volume': [5200, 8100, 3400]  # 设置仅用于类型转换演示的假设成交量
})  # 完成含字符串、浮点和整数列的异构示例表
arrow_table = pa.Table.from_pandas(sample_stock_df, preserve_index=False)  # 转为 Arrow 表;字符串编码和缓冲区适配可能发生复制
print(f'Arrow 表的 Schema:\n{arrow_table.schema}')  # 输出 Arrow 表的字段类型信息
Arrow 表的 Schema:
stock_code: string
close_price: double
volume: int64
-- schema metadata --
pandas: '{"index_columns": [], "column_indexes": [], "columns": [{"name":' + 445
recovered_df = arrow_table.to_pandas()  # 转回 Pandas;默认行为不承诺与 Arrow 缓冲区零拷贝共享
print(f'转换后的 DataFrame:\n{recovered_df}')  # 验证数据完整性
print(f'\n数据类型保持一致: {recovered_df.dtypes.tolist()}')  # 确认数据类型
转换后的 DataFrame:
  stock_code  close_price  volume
0     600276         45.0    5200
1     002415         32.0    8100
2     600104         15.0    3400

数据类型保持一致: [dtype('O'), dtype('float64'), dtype('int64')]

4.6 格式选择指南

根据数据规模和使用场景,推荐以下存储格式:

场景 推荐格式 理由
小型数据(< 100 MB),需要人工检查 CSV 通用性最好,可用文本编辑器查看
需要键管理,且节点可写为 table/data_columns HDF5 满足布局前提时可用 where 条件筛选
需要跨引擎扫描、列投影或分区过滤 Parquet 是否更小或更快取决于编码、压缩、列基数与查询
兼容列式缓冲区的进程内交换 Arrow 类型与缓冲区兼容时可能减少复制
跨进程或持久化 Arrow 消息 Arrow IPC 保留 Arrow schema,但仍有消息序列化与 I/O 成本
数据库查询场景 SQLite / DuckDB 支持 SQL,适合复杂聚合查询

4.7 本章小结

本章介绍了数据分析工作流中至关重要的数据 I/O 环节:

  1. CSV 进阶参数parse_datesdtypeusecols 等参数可以大幅提升读取效率和减少内存占用
  2. 分块迭代chunksize 参数使得 Pandas 能够处理超出内存的大型 CSV 文件
  3. SQL 数据库:通过 SQLAlchemysqlite3,Pandas 可以直接读写关系型数据库
  4. 列式文件与内存:Parquet 以 row group、column chunk 和 page 组织文件,Arrow 定义内存数组与缓冲区;实际速度、体积与复制边界需在具体 schema 和查询上验证
  5. HDF5 选择性读取where 只适用于 table 节点及其可查询列;是否全量读取应由分析范围、内存预算与后续算子共同决定

4.8 练习题

练习 4.1:使用本地数据集 stock_price_pre_adjusted.h5,先按证券代码选择性读取海康威视全部行情列,再在相同条件下仅投影 closevolume 两列;比较耗时和内存占用。为什么教材不要求全量载入整个市场文件?

解答

import pandas as pd  # 导入 Pandas
import time  # 导入时间模块用于计时
stock_price_path = f'{DATA_ROOT}/stock/stock_price_pre_adjusted.h5'  # 构建文件路径
start_all_columns = time.perf_counter()  # 记录读取同一公司全部列的起始时刻
hikvision_all_columns = pd.read_hdf(  # 在安全的单公司子集上读取全部行情列
    stock_price_path, where="order_book_id == '002415.XSHE'"  # 固定证券代码以保证两方案行集一致
)  # 完成全部列方案的选择性读取
all_columns_seconds = time.perf_counter() - start_all_columns  # 计算本机本次全部列读取耗时

start_projection = time.perf_counter()  # 记录同一公司列投影方案的起始时刻
hikvision_selective = pd.read_hdf(  # 同时执行证券行过滤和所需列投影
    stock_price_path,  # 复用完全相同的 HDF5 文件
    where="order_book_id == '002415.XSHE'",  # 保持与全部列方案相同的行条件
    columns=['close', 'volume']  # 只物化收盘价和成交量两列
)  # 完成列投影方案的选择性读取
projection_seconds = time.perf_counter() - start_projection  # 计算本机本次投影读取耗时

comparison_df = pd.DataFrame({  # 汇总同一行集下两种读取方案的观测指标
    '读取方案': ['行筛选:全部列', '行筛选+列投影'],  # 标明每行对应的物化列范围
    '耗时_秒': [all_columns_seconds, projection_seconds],  # 保存本机单次计时而非普遍性能结论
    '内存_MB': [  # 计算两个结果对象各自的深度内存占用
        hikvision_all_columns.memory_usage(deep=True).sum() / 2**20,  # 全部列结果的内存量
        hikvision_selective.memory_usage(deep=True).sum() / 2**20  # 两列投影结果的内存量
    ]  # 完成内存指标列表
})  # 完成两方案比较表
comparison_df  # 展示计时与内存结果,解释时保留缓存和硬件条件
读取方案 耗时_秒 内存_MB
0 行筛选:全部列 5.811578 0.339487
1 行筛选+列投影 8.802307 0.223856

完整市场文件远大于单只股票子集,全量读取既违背本项目的选择性载入规范,也可能耗尽课堂机器内存。这里比较的是两个安全方案;结果会随磁盘缓存而波动,因此应重复测量并重点解释内存差异,而不是把一次计时当成普遍速度结论。

练习 4.2:将练习 4.1 中筛选出的海康威视股价数据分别保存为 CSV、Parquet 和 HDF5 三种格式,比较文件大小差异。

解答

import os  # 导入os模块用于获取文件大小
hikvision_csv_path = chapter_io_root / 'hikvision_price.csv'  # 在本章临时目录指定CSV比较文件
hikvision_parquet_path = chapter_io_root / 'hikvision_price.parquet'  # 在同一临时目录指定Parquet比较文件
hikvision_hdf5_path = chapter_io_root / 'hikvision_price.h5'  # 在同一临时目录指定HDF5比较文件
hikvision_selective.to_csv(hikvision_csv_path, index=False)  # 把CSV比较文件写入临时目录
hikvision_selective.to_parquet(hikvision_parquet_path)  # 把Parquet比较文件写入临时目录
hikvision_selective.to_hdf(hikvision_hdf5_path, key='price', mode='w')  # 把HDF5比较文件写入临时目录
csv_size = os.path.getsize(hikvision_csv_path) / 1024  # 获取临时CSV文件大小(KB)
parquet_size = os.path.getsize(hikvision_parquet_path) / 1024  # 获取临时Parquet文件大小(KB)
hdf5_size = os.path.getsize(hikvision_hdf5_path) / 1024  # 获取临时HDF5文件大小(KB)
print(f'CSV:     {csv_size:.1f} KB')  # 输出CSV大小
print(f'Parquet: {parquet_size:.1f} KB')  # 输出Parquet大小
print(f'HDF5:    {hdf5_size:.1f} KB')  # 输出HDF5大小
CSV:     68.5 KB
Parquet: 90.8 KB
HDF5:    109.0 KB

文件大小不是格式的固定属性:CSV 的文本表示、Parquet 的编码和压缩、HDF5 的布局和压缩设置都会改变结果。本题只比较由同一 DataFrame 在当前默认参数下写出的三个文件。

练习 4.3:把本章 SQL 查询改写为 Pandas groupby(),核对每个证券代码的记录数、平均价格与总成交量是否一致。为什么只比较行数不足以证明两个结果相同?

解答

pandas_summary = (  # 用 Pandas 重现 SQL 的分组聚合语义
    trade_data_for_db.groupby('stock_code', observed=True)  # 按证券代码分组并避免生成未观察分类
    .agg(trade_count=('price', 'size'), avg_price=('price', 'mean'), total_volume=('volume', 'sum'))  # 计算三项同口径指标
    .sort_values('total_volume', ascending=False)  # 对齐 SQL 的成交量降序排序规则
    .reset_index()  # 把分组键恢复为普通列以便逐列比较
)  # 完成 Pandas 聚合结果
sql_check_df = stock_summary_from_sql.sort_values('stock_code').reset_index(drop=True)  # 按键排序 SQL 结果以消除展示顺序差异
pandas_check_df = pandas_summary.sort_values('stock_code').reset_index(drop=True)  # 按相同键排序 Pandas 结果
pd.testing.assert_frame_equal(sql_check_df, pandas_check_df, check_dtype=False)  # 同时核对键、全部指标和数值,而非只核对行数
pandas_summary  # 展示已通过一致性断言的 Pandas 聚合结果
stock_code trade_count avg_price total_volume
0 002415.XSHE 3789 21.505153 1.468341e+11

练习 4.4:判断下列说法是否成立并说明条件:“Pandas 转 Arrow 一定零拷贝”“Arrow IPC 没有序列化开销”“任意 HDF5 节点都能使用 where”。

解答:三个说法都不成立。Pandas 与 Arrow 只有在 dtype、缓冲区布局、分块和所有权兼容时才可能共享缓冲区,字符串或对象列常需编码与分配;IPC 要组织 schema、元数据和缓冲区并执行 I/O;Pandas 的 HDF5 where 查询要求 table 布局,且筛选字段必须是可查询的 data_columns。实际工作应先检查 schema 和节点元数据,再用小规模往返与内存共享测试验证假设。

练习 4.5:对本章在临时目录生成的 demo_trades.csv 只读取 stock_codepricevolume 三列,显式指定 dtype,并以 chunksize=500 分块计算行数、价格总和和成交量总和。再安全地全量读取这份单公司教学副本,逐项核对分块结果与全量结果;说明为什么只看到代码执行完成不足以证明分块逻辑正确。

解答

import numpy as np  # 使用数值容差核对分块与全量聚合,避免脆弱的浮点精确比较
csv_usecols = ['stock_code', 'price', 'volume']  # 只物化完成聚合所需字段以检验列投影
csv_dtypes = {'stock_code': 'string', 'price': 'float64', 'volume': 'float64'}  # 显式固定代码与数值列的读取类型
full_csv_subset = pd.read_csv(demo_csv_path, usecols=csv_usecols, dtype=csv_dtypes)  # 从同一临时路径全量读取安全单公司副本作为基准
assert full_csv_subset.columns.tolist() == csv_usecols  # 核对usecols结果没有遗漏或意外增加字段
assert full_csv_subset.dtypes.astype(str).to_dict() == csv_dtypes  # 核对实际dtype与声明的数据合同一致
chunk_row_count = 0  # 初始化分块累计行数以检测漏块或重复处理
chunk_price_sum = 0.0  # 初始化价格总和作为数值正确性检查量
chunk_volume_sum = 0.0  # 初始化成交量总和作为业务聚合检查量
chunk_reader = pd.read_csv(demo_csv_path, usecols=csv_usecols, dtype=csv_dtypes, chunksize=500)  # 从同一临时路径用相同schema创建分块迭代器
for trade_chunk in chunk_reader:  # 逐块消费文件而不同时保留所有分块对象
    chunk_row_count += len(trade_chunk)  # 累加当前分块行数以覆盖完整文件
    chunk_price_sum += trade_chunk['price'].sum()  # 累加价格列以验证数值读取和块边界
    chunk_volume_sum += trade_chunk['volume'].sum()  # 累加成交量以验证核心业务汇总
full_totals = full_csv_subset[['price', 'volume']].sum()  # 在安全全量基准上独立计算两项总和
assert chunk_row_count == len(full_csv_subset)  # 核对分块路径既没有漏行也没有重复行
np.testing.assert_allclose([chunk_price_sum, chunk_volume_sum], full_totals.to_numpy())  # 在浮点容差内核对全部聚合值
print(f'核验通过:{chunk_row_count:,} 行,价格和={chunk_price_sum:.2f},成交量和={chunk_volume_sum:.0f}')  # 报告已通过对照的分块结果
核验通过:3,789 行,价格和=81483.03,成交量和=146834111964

程序“不报错”只能说明语法与运行路径成立,不能排除漏掉末块、重复累计或读错列。行数与两个独立数值汇总同时一致,才为这项分块任务提供了可观察的正确性证据;换到无法安全全量读取的大文件时,应先抽取可复核子集完成同类测试。