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教学副本的临时路径4 数据读取与存储:从本地到大数据系统
4.1 引言与学习目标
数据分析的第一步是将数据从各种来源读入 Python 环境。在真实的金融数据分析工作中,数据可能存储在 CSV 文件、HDF5 文件、Parquet 文件、SQL 数据库,甚至是网络 API 中。选择合适的数据格式和读写方法,不仅影响代码的运行效率,还直接决定了能否在有限的内存中处理大规模数据集。
学习目标
完成本章后,你应能:
- 根据按行读取、列投影、谓词过滤、跨工具交换和人工检查需求,比较 CSV、HDF5、Parquet 与 Arrow 的适用边界;
- 识别
DataFrame的行、列与索引,完成列选择、布尔行筛选、索引设置以及groupby().agg()最小聚合; - 使用
read_csv()的dtype、usecols和chunksize,并用全量安全子集核对分块聚合结果; - 创建 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 及答案中的选择、聚合、排序与逐项核对 |
正确使用 dtype、usecols 与 chunksize |
小节 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 小节从 SELECT、GROUP BY 和聚合函数的最小例开始,不要求数据库先修。
4.2 结构化文本读写(CSV)的进阶参数与分块迭代
4.2.1 CSV 读取的核心参数
pd.read_csv() 是 Pandas 中最常用的数据读取函数。除了基本的文件路径参数外,掌握以下进阶参数对于处理真实世界的数据至关重要:
教师提供的数据准备(只运行核对,不评分)
下一块用尚未正式桥接的 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-8 或 gbk |
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 展示本节讨论对象的可视化结果,读图时应结合正文给出的口径与限制。
4.5.2 HDF5 格式深度解析
HDF5(Hierarchical Data Format 5)是一种支持分层数据组织的文件格式,特别适合存储大型科学数据集。本教材的核心数据集均采用 HDF5 格式。
HDF5 的核心优势
- 分层结构:类似文件系统,一个 HDF5 文件内部可以包含多个”组”和”数据集”
- 选择性读取有前提:Pandas 只有在节点以
format='table'写入,且筛选列被声明为data_columns时,才能用where下推相应条件;fixed节点不支持这种查询 - 压缩支持:支持多种压缩算法(如 blosc、zlib),大幅减小文件体积
- 类型丰富:原生支持多维数组、字符串、日期时间等类型
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 环节:
- CSV 进阶参数:
parse_dates、dtype、usecols等参数可以大幅提升读取效率和减少内存占用 - 分块迭代:
chunksize参数使得 Pandas 能够处理超出内存的大型 CSV 文件 - SQL 数据库:通过
SQLAlchemy和sqlite3,Pandas 可以直接读写关系型数据库 - 列式文件与内存:Parquet 以 row group、column chunk 和 page 组织文件,Arrow 定义内存数组与缓冲区;实际速度、体积与复制边界需在具体 schema 和查询上验证
- HDF5 选择性读取:
where只适用于 table 节点及其可查询列;是否全量读取应由分析范围、内存预算与后续算子共同决定
4.8 练习题
练习 4.1:使用本地数据集 stock_price_pre_adjusted.h5,先按证券代码选择性读取海康威视全部行情列,再在相同条件下仅投影 close 与 volume 两列;比较耗时和内存占用。为什么教材不要求全量载入整个市场文件?
解答:
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_code、price 与 volume 三列,显式指定 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
程序“不报错”只能说明语法与运行路径成立,不能排除漏掉末块、重复累计或读错列。行数与两个独立数值汇总同时一致,才为这项分块任务提供了可观察的正确性证据;换到无法安全全量读取的大文件时,应先抽取可复核子集完成同类测试。