先修知识与学习目标
本章承接 Pandas 的索引、分组聚合和时间序列知识。读者应会使用 read_hdf()、groupby() 和 astype(),并理解整数、浮点数与字符串的基本类型。这里讨论的核心问题不是笼统地问 Pandas 能处理多少行,而是判断一次具体计算的峰值内存 、数据复制和算法复杂度。
完成本章后,读者应能:
目标 11.1 :解释 NumPy 缓冲区、Python 对象、索引和缺失值掩码如何共同形成 DataFrame 的内存占用;
目标 11.2 :根据取值范围与误差容忍度选择整数、浮点、分类和 Arrow 类型;
目标 11.3 :使用选择性读取与分块聚合控制峰值内存,并验证分块结果;
目标 11.4 :识别降级溢出、浮点精度损失、类别基数过高和链式计算隐式复制等风险。
目标—活动—核心练习—答案证据映射
表 表 11.1 汇总本节的计算或审计结果,解释时应遵循正文给出的口径与限制。
内存模型、复杂度与边界
一个 DataFrame 的内存并不等于文件大小。忽略分配器碎片时,可用 式 11.1 近似表示:
\[
M_{\text{frame}} \approx M_{\text{index}} + \sum_{j=1}^{p}\left(M_{\text{buffer},j}+M_{\text{mask},j}+M_{\text{object},j}\right).
\tag{11.1}\]
定长数值列主要由连续缓冲区构成;可空类型还需要有效性掩码。传统 object 字符串列的数组保存 Python 对象指针,字符串内容另在堆中分配,因此浅层统计会低估成本。memory_usage(deep=True) 会递归估计对象内容,更适合教学中的相对比较,但它仍不是进程峰值内存的精确测量。
若一列有 \(n\) 个值、\(k\) 个不同类别,类别编码宽度为 \(b_c\) 字节,类别字典平均每项成本为 \(b_d\) ,则分类存储成本近似为:
\[
M_{\text{category}} \approx n b_c + k b_d.
\tag{11.2}\]
当 \(k \ll n\) 时,式 11.2 通常远小于逐行保存字符串;当几乎每行都不同,类别字典反而可能增加内存。类型转换需要扫描全部 \(n\) 个值,时间复杂度通常为 \(O(n)\) ,并可能在转换期间同时保留新旧缓冲区,使峰值接近转换前后内存之和。
表 表 11.2 汇总本节的计算或审计结果,解释时应遵循正文给出的口径与限制。
其中 \(c\) 是块大小,\(g\) 是分组数。分块并不自动保证正确:求和、计数、最小值和最大值可直接合并;均值必须同时维护分子与分母;中位数和全局排序则不能仅合并各块结果。
本地真实行情的选择性读取
本章使用上海、江苏、浙江的上汽集团、恒瑞医药、海康威视和宁波银行 2023 年前复权日行情。第一段代码先建立跨平台路径,再只读取所需证券、期间和列;这比先载入 1536 万行全表再筛选更符合内存约束。
import platform # 识别操作系统以复用同一份教材代码
from pathlib import Path # 以结构化路径表达本地数据位置
import numpy as np # 提供数值类型边界与精度检查工具
import pandas as pd # 执行HDF选择性读取和内存分析
import pyarrow.dataset as ds # 以记录批次扫描真实Parquet分片
DATA_ROOT = 'C:/qiufei/data' if platform.system() == 'Windows' else '/home/ubuntu/r2_data_mount/data' # 选择规范数据根目录
PRICE_PATH = Path(DATA_ROOT) / 'stock' / 'stock_price_pre_adjusted.h5' # 指向前复权日行情文件
MARKET_DATASET = Path(DATA_ROOT) / 'stock' / 'full_market_with_return_2023.parquet' # 指向可批量扫描的真实Parquet目录
YRD_CODES = ['600104.XSHG' , '600276.XSHG' , '002415.XSHE' , '002142.XSHE' ] # 固定长三角公司样本
PRICE_COLUMNS = ['open' , 'close' , 'volume' , 'total_turnover' ] # 仅请求分析所需字段
PRICE_FILTERS = [f'order_book_id in { YRD_CODES} ' , 'date >= Timestamp( \' 2023-01-01 \' )' , 'date <= Timestamp( \' 2023-12-31 \' )' ] # 在HDF层过滤证券与日期
stock_prices = pd.read_hdf(PRICE_PATH, key= 'data' , where= PRICE_FILTERS, columns= PRICE_COLUMNS).reset_index() # 只把目标切片载入内存
assert set (stock_prices.columns) == {'order_book_id' , 'date' , 'open' , 'close' , 'volume' , 'total_turnover' } # 显式核验真实schema
print (stock_prices.shape) # 报告实际载入规模
print (stock_prices.dtypes) # 展示HDF读取后的原始类型
(968, 6)
order_book_id object
date datetime64[ns]
open float64
close float64
volume float64
total_turnover float64
dtype: object
表 11.3 报告真实切片的深层内存,而不是用随机行情代替市场数据。
memory_profile = pd.DataFrame({'dtype' : stock_prices.dtypes.astype(str ), 'minimum' : stock_prices.select_dtypes('number' ).min (), 'maximum' : stock_prices.select_dtypes('number' ).max (), 'memory_kb' : stock_prices.memory_usage(deep= True ) / 1024 }) # 合并类型、范围和深层内存证据
memory_profile.index.name = 'column' # 明确表格索引的业务含义
memory_profile.round (3 ) # 输出便于阅读的诊断表
安全的类型优化
整数降级:先证明范围,再转换
把成交量机械地改成 int16 会溢出。正确顺序是先检查最小值、最大值和缺失值,再选能覆盖样本与合理未来范围的类型。这里成交量非负,因此无缺失时可以使用无符号整数。
has_missing_volume = stock_prices['volume' ].isna().any () # 判断普通NumPy整数是否能表达该列
maximum_volume = stock_prices['volume' ].max () # 提取样本最大成交量用于边界证明
assert not has_missing_volume # 确认转换不会丢失缺失值语义
assert maximum_volume <= np.iinfo(np.uint32).max # 确认所有观测落在uint32范围内
optimized_prices = stock_prices.copy() # 保留基准版本以比较结果和内存
optimized_prices['volume' ] = optimized_prices['volume' ].astype('uint32' ) # 在边界验证后压缩成交量
print (optimized_prices['volume' ].dtype) # 输出最终成交量类型
浮点降级:误差容忍度是业务约束
float32 只有约七位十进制有效数字。价格降级是否安全,不能由数据类型名称决定,而应由计算用途决定。若教学任务允许最大绝对价格误差不超过 \(10^{-3}\) 元,可用 式 11.3 验证:
\[
e_{\max}=\max_i\left|x_i-\operatorname{float32}(x_i)\right|.
\tag{11.3}\]
close_float32 = optimized_prices['close' ].astype('float32' ) # 生成候选低精度价格缓冲区
maximum_price_error = (optimized_prices['close' ] - close_float32.astype('float64' )).abs ().max () # 在同一精度下计算绝对误差
assert maximum_price_error < 0.001 # 按千分之一元容忍度判断本案例是否可接受
optimized_prices['open' ] = optimized_prices['open' ].astype('float32' ) # 在相同价格量级下压缩开盘价
optimized_prices['close' ] = close_float32 # 应用已验证的收盘价类型
print (f'最大价格误差为 { maximum_price_error:.8f} 元' ) # 报告精度损失而非宣称无损
在逐笔成交金额、会计金额核对或长期复利计算中,误差可能累积;此时应保留 float64,或使用定点整数表示最小货币单位。降级是有条件的工程决策,不是越小越好。
分类类型:基数决定收益
股票代码只有四个不同值,是低基数列。转换前后必须同时比较内存与值一致性。
code_memory_before_kb = optimized_prices['order_book_id' ].memory_usage(deep= True ) / 1024 # 度量对象字符串的深层内存
original_codes = optimized_prices['order_book_id' ].copy() # 保存原值用于无损性验证
optimized_prices['order_book_id' ] = optimized_prices['order_book_id' ].astype('category' ) # 将低基数证券代码编码为分类类型
code_memory_after_kb = optimized_prices['order_book_id' ].memory_usage(deep= True ) / 1024 # 度量编码和字典的总内存
assert original_codes.equals(optimized_prices['order_book_id' ].astype('object' )) # 验证编码前后业务值完全一致
print (f'代码列内存由 { code_memory_before_kb:.2f} KB 降至 { code_memory_after_kb:.2f} KB' ) # 报告本机实测结果
代码列内存由 64.41 KB 降至 1.50 KB
Arrow 后端:列式布局不等于无条件零拷贝
Arrow 字符串把偏移量、字节内容和有效性位图组织为列式缓冲区,避免逐值 Python 对象。Pandas 可用 convert_dtypes(dtype_backend='pyarrow') 获得 Arrow 扩展类型,但转换已有 NumPy/object 列通常需要分配新缓冲区。后续与 Arrow 生态互操作能否零拷贝,还取决于类型是否兼容、缓冲区是否连续、是否需要合并 chunks,以及接收方是否接受只读内存。
arrow_prices = stock_prices.convert_dtypes(dtype_backend= 'pyarrow' ) # 将兼容列转换为Pandas Arrow扩展类型
arrow_memory_kb = arrow_prices.memory_usage(deep= True ).sum () / 1024 # 度量转换后的DataFrame内存
original_memory_kb = stock_prices.memory_usage(deep= True ).sum () / 1024 # 度量同一观测集的原始内存
print (arrow_prices.dtypes) # 检查每列实际采用的Arrow类型
print (f'原始内存 { original_memory_kb:.2f} KB,Arrow类型内存 { arrow_memory_kb:.2f} KB' ) # 报告特定环境的实测值
order_book_id string[pyarrow]
date timestamp[ns][pyarrow]
open double[pyarrow]
close double[pyarrow]
volume int64[pyarrow]
total_turnover double[pyarrow]
dtype: object
原始内存 102.22 KB,Arrow类型内存 52.12 KB
易混淆概念辨析:列式存储、Arrow 后端与零拷贝
列式存储描述同一列值的组织方式;Arrow 是一种内存格式;零拷贝则是一次具体交换是否共享既有缓冲区。三者相关但不等价。read_hdf() 读取 HDF5 后再转换为 Arrow 已经发生了解码与分配,因此不能把整条链路称为零拷贝。
共享视图、Copy-on-Write 与隐式复制
“没有写 .copy()”不等于没有复制。相同 NumPy 类型的 astype(..., copy=False) 可能共享既有缓冲区;类型降级必须新建另一种元素宽度的缓冲区;算术表达式、布尔筛选、排序和许多跨库转换也通常产生结果数组。Pandas 的 Copy-on-Write 主要约束对象修改时的隔离语义,不会消除计算本身所需的临时数组。
因此,内存审计至少回答两个问题:最终对象占多少内存,以及转换期间是否同时存活源缓冲区、目标缓冲区和中间结果。np.shares_memory() 可以为 NumPy 支持的数组提供局部证据,但不能据此推断整条 DataFrame 或 Arrow 链路都零拷贝。小节 11.8.5 将类型不变、类型降级和算术临时结果放在同一审计表中。
分块计算与峰值内存
若最终结果只是每只股票的总成交额,不必保留全部明细。对可加总指标,分块算法维护 \(g\) 个组的部分和,工作内存由全量 \(O(n)\) 降为约 \(O(c+g)\) 。
HDF5 已在前文用于选择性读取。为了让批次边界由文件扫描器真正控制,下面对同一年度的本地 Parquet 分片使用 Arrow record batch;这避免 PyTables 对大表迭代筛选时的额外扫描状态。
turnover_parts = [] # 收集每个块的证券级部分和
market_dataset = ds.dataset(MARKET_DATASET, format = 'parquet' ) # 从真实Parquet分片建立只读数据集
market_filter = ds.field('order_book_id' ).isin(YRD_CODES) & (ds.field('trade_date' ) < pd.Timestamp('2024-01-01' )) # 把证券和年度筛选下推到Arrow扫描器
market_scanner = market_dataset.scanner(columns= ['order_book_id' , 'total_turnover' ], filter = market_filter, batch_size= 100_000 ) # 将物理扫描批次限制为至多十万行
for record_batch in market_scanner.to_batches(): # 逐批推进而不拼接日度明细
price_chunk = record_batch.to_pandas() # 只把当前小批次转换为Pandas
chunk_totals = price_chunk.groupby('order_book_id' )['total_turnover' ].sum () # 在块内先压缩为证券级统计量
turnover_parts.append(chunk_totals) # 保存小型中间结果而非行情明细
chunked_turnover = pd.concat(turnover_parts, axis= 1 ).fillna(0 ).sum (axis= 1 ).sort_values(ascending= False ) # 合并各块的可加总统计量
direct_parquet_prices = market_dataset.to_table(columns= ['order_book_id' , 'total_turnover' ], filter = market_filter).to_pandas() # 仅为验证物化同源的小型目标切片
direct_turnover = direct_parquet_prices.groupby('order_book_id' )['total_turnover' ].sum ().sort_values(ascending= False ) # 计算同一Parquet来源的内存内基准结果
pd.testing.assert_series_equal(chunked_turnover, direct_turnover, check_names= False , check_index_type= False ) # 忽略跨格式字符串索引实现并核验数值一致
chunked_turnover.to_frame('total_turnover_yuan' ) # 输出各公司全年成交额
order_book_id
600276.XSHG
3.781830e+11
002415.XSHE
3.286586e+11
002142.XSHE
1.993612e+11
600104.XSHG
6.962311e+10
结果解释
以上输出是本地真实行情在当前 peter 环境中的测量值。分类列收益较大,是因为四个代码被重复使用约千次;数值降级收益较温和,因为日期和成交额仍保留原类型。分块聚合与直接聚合的一致性断言说明:对求和而言,块内先聚合再合并没有改变结果。它并不证明所有分析都可分块,尤其不能直接推广到精确中位数、全局排名或跨块滚动窗口。
常见误区
只看文件大小估计内存。 压缩 Parquet/HDF5 解码后可能显著膨胀,操作中的临时数组还会提高峰值。
把所有整数改成最小类型。 新数据超出历史范围会溢出;含缺失值时普通 NumPy 整数还会丢失缺失语义。
认为 float32 对金融价格天然安全。 单价展示与组合净值累计的误差要求不同,必须用业务容忍度验证。
把所有字符串改成 category。 高基数字段会维护庞大字典,频繁拼接不同类别集合也有重编码成本。
把分块读取等同于超内存计算。 若把每个块追加到列表后再次拼成全表,峰值内存并未真正受控。
本章小结
Pandas 的边界由数据布局、算子和峰值内存共同决定,而不是由固定行数决定。安全优化遵循同一审计链:先选择性读取,再测量深层内存与取值范围,随后应用类型或分块策略,最后验证数值和业务语义不变。Arrow 能改善列式互操作,但零拷贝必须逐条链路判断。
分层练习与完整答案
练习 11.1:概念检查
某证券代码列有 \(n=10^7\) 行、\(k=5000\) 个类别。类别编码使用 2 字节,类别字典平均每项 80 字节。忽略其他开销,估算分类存储内存,并解释为何不能仅凭该结果保证转换时不会内存溢出。
答案
由 式 11.2 ,\(M=10^7\times2+5000\times80=20{,}400{,}000\) 字节,约 19.46 MiB。转换阶段通常同时存在原 object 缓冲区、Python 字符串对象、新编码和类别字典,因此峰值可能远大于最终 19.46 MiB;还要计入索引、掩码和分配器开销。
练习 11.2:类型边界诊断
验证本章真实样本的成交量能否使用 uint16,并给出正确选择。
答案
sample_maximum_volume = stock_prices['volume' ].max () # 提取真实样本的最大成交量
uint16_maximum = np.iinfo(np.uint16).max # 取得uint16可表达的上界
can_use_uint16 = sample_maximum_volume <= uint16_maximum # 对照数据范围与类型上界
recommended_dtype = 'uint16' if can_use_uint16 else 'uint32' # 选择能够覆盖观测的较小无符号类型
print (f'样本最大成交量为 { sample_maximum_volume:.0f} ,uint16上界为 { uint16_maximum} ' ) # 输出判断证据
print (f'推荐类型为 { recommended_dtype} ' ) # 输出可复核的类型结论
样本最大成交量为 137880485,uint16上界为 65535
推荐类型为 uint32
输出应显示 uint16 不足,因为 A 股日成交量可以远高于 65,535 股;当前样本应选择 uint32。
练习 11.3:真实基础信息优化
读取 stock_basic_data.h5,使用经核验的真实字段 province、industry_name 和 board_type,比较转为分类类型前后的内存,并确保值不变。
答案
BASIC_PATH = Path(DATA_ROOT) / 'stock' / 'stock_basic_data.h5' # 指向真实股票基础信息文件
stock_basic = pd.read_hdf(BASIC_PATH, key= 'stock_basic_info' ) # 读取规模较小的完整基础信息表
category_columns = ['province' , 'industry_name' , 'board_type' ] # 使用实测存在的地区、行业与板块字段
assert set (category_columns).issubset(stock_basic.columns) # 防止错误schema被静默跳过
basic_values_before = stock_basic[category_columns].copy() # 保存业务值作为转换基准
basic_memory_before = stock_basic.memory_usage(deep= True ).sum () # 计算转换前深层内存
stock_basic[category_columns] = stock_basic[category_columns].astype('category' ) # 编码三列低基数字段
basic_memory_after = stock_basic.memory_usage(deep= True ).sum () # 计算转换后深层内存
pd.testing.assert_frame_equal(basic_values_before, stock_basic[category_columns].astype('object' )) # 验证所有文本值保持一致
print (f'总内存由 { basic_memory_before / 1024 ** 2 :.2f} MiB 降至 { basic_memory_after / 1024 ** 2 :.2f} MiB' ) # 报告本地实测结果
总内存由 8.22 MiB 降至 6.93 MiB
结果应解释为该文件与环境下的观测值,而不是可推广到任意数据的固定节省比例。
练习 11.4:综合挑战
用分块方式同时计算四只股票的成交额总和与记录数,再据此计算平均日成交额;验证它等于直接聚合结果。说明为何不能对各块均值再做简单平均。
答案
chunk_statistics = [] # 保存每块的求和与计数充分统计量
exercise_scanner = market_dataset.scanner(columns= ['order_book_id' , 'total_turnover' ], filter = market_filter, batch_size= 100_000 ) # 为练习创建新的真实行情批次扫描器
for exercise_batch in exercise_scanner.to_batches(): # 逐批读取目标证券观测
exercise_chunk = exercise_batch.to_pandas() # 只物化当前批次
grouped_chunk = exercise_chunk.groupby('order_book_id' )['total_turnover' ].agg(['sum' , 'count' ]) # 为均值保留分子和分母
chunk_statistics.append(grouped_chunk) # 累积小型分组统计而非原始行
combined_statistics = pd.concat(chunk_statistics).groupby(level= 0 ).sum () # 跨块分别合并总额与记录数
combined_statistics['mean' ] = combined_statistics['sum' ] / combined_statistics['count' ] # 用全局分子除以全局分母
combined_statistics # 输出分块均值及其组成
order_book_id
002142.XSHE
1.993612e+11
241
8.272249e+08
002415.XSHE
3.286586e+11
241
1.363729e+09
600104.XSHG
6.962311e+10
241
2.888926e+08
600276.XSHG
3.781830e+11
241
1.569224e+09
direct_mean = direct_parquet_prices.groupby('order_book_id' )['total_turnover' ].mean().sort_index() # 计算同一Parquet来源的直接验证基准
chunk_mean = combined_statistics['mean' ].sort_index() # 对齐分块结果的证券顺序
pd.testing.assert_series_equal(chunk_mean, direct_mean, check_names= False , check_index_type= False ) # 忽略跨格式字符串索引实现并验证数值一致
print (chunk_mean) # 输出最终平均日成交额
order_book_id
002142.XSHE 8.272249e+08
002415.XSHE 1.363729e+09
600104.XSHG 2.888926e+08
600276.XSHG 1.569224e+09
Name: mean, dtype: float64
各块行数可能不同,简单平均各块均值会给小块与大块相同权重;保留 sum 和 count 才能得到正确的加权结果。
练习 11.5:Arrow 类型、浮点精度与隐式复制审计
针对本章真实行情回答三组问题:
分别为 order_book_id、date、volume 和 close 选择适合“本章 Pandas 分析”或“交给 Arrow 引擎”的类型,并说明选择条件;
比较 float64 与 float32 日收益复利财富路径的最大绝对差,判断为什么单价展示通过 \(10^{-3}\) 元阈值不等于长期复利或金额核对也安全;
用共享内存检查区分相同类型 astype(copy=False)、float32 降级与乘法表达式,指出哪两步必然需要新数组。
评分点(100 分) :四列类型选择与条件 35 分;精度计算和业务判断 30 分;三种内存共享结果 25 分;不夸大零拷贝结论 10 分。
答案
表 表 11.4 汇总本节的计算或审计结果,解释时应遵循正文给出的口径与限制。
表 表 11.5 汇总本节的计算或审计结果,解释时应遵循正文给出的口径与限制。
anchor_close_float64 = stock_prices.loc[stock_prices['order_book_id' ].eq('600276.XSHG' )].sort_values('date' )['close' ].astype('float64' ) # 取恒瑞医药真实价格作为共同输入
anchor_returns_float64 = anchor_close_float64.pct_change(fill_method= None ).dropna() # 以双精度计算相邻真实交易日收益
wealth_float64 = (1 + anchor_returns_float64).cumprod() # 形成双精度财富路径基准
wealth_float32 = (1 + anchor_returns_float64.astype('float32' )).cumprod().astype('float64' ) # 在单精度收益上执行同一复利递推
maximum_wealth_path_gap = (wealth_float64 - wealth_float32).abs ().max () # 度量舍入在多期乘法中的累计影响
same_dtype_view = anchor_close_float64.astype('float64' , copy= False ) # 请求在类型不变时复用源缓冲区
narrowed_float32 = anchor_close_float64.astype('float32' ) # 元素宽度变化要求分配新的目标缓冲区
scaled_close = anchor_close_float64 * 100 # 算术表达式生成新的结果数组而非原地改写价格
shares_same_dtype = np.shares_memory(anchor_close_float64.to_numpy(copy= False ), same_dtype_view.to_numpy(copy= False )) # 检查类型不变路径是否实际共享内存
shares_narrowed = np.shares_memory(anchor_close_float64.to_numpy(copy= False ), narrowed_float32.to_numpy(copy= False )) # 检查降级数组是否错误复用不同宽度缓冲区
shares_scaled = np.shares_memory(anchor_close_float64.to_numpy(copy= False ), scaled_close.to_numpy(copy= False )) # 检查乘法结果是否另行分配
assert shares_same_dtype and not shares_narrowed and not shares_scaled # 固定三种路径的可评分内存合同
pd.DataFrame({'audit' : ['float32复利路径最大绝对差' , '同类型astype共享内存' , 'float32降级共享内存' , '乘法结果共享内存' ], 'observed' : [maximum_wealth_path_gap, shares_same_dtype, shares_narrowed, shares_scaled]}) # 输出精度与复制证据
答案中的 maximum_wealth_path_gap 必须如实报告而不能预设为零。即便本样本的误差很小,也只能支持这条财富路径在当前容忍度下可接受;会计对账、逐笔金额或更长持有期需要重新定义阈值。后两项 False 说明类型降级和乘法表达式产生新数组;第一项 True 只是这个同类型 NumPy 缓冲区的局部证据,不能推广成整个 DataFrame 或 Arrow 流程无条件零拷贝。