12  Polars高性能DataFrame

12.1 先修知识与学习目标

本章假定读者会使用 Pandas 的列选择、布尔筛选、分组聚合与窗口计算,并理解第 11 章的列式内存概念。Polars 不是 Pandas 的无缝加速开关:它没有 Pandas 的隐式行索引,表达式语义、缺失值处理和执行计划也不同。

完成本章后,读者应能:

  • 用表达式描述列变换,并解释表达式为何便于并行与查询优化;
  • 区分即时执行、延迟执行与流式执行,识别三者各自的物化边界;
  • 从本地真实 Parquet 数据构建带谓词下推和投影下推的查询;
  • 在固定数据、查询、环境和计时口径下比较 Polars 与 Pandas,而不使用无依据的性能倍数。

本章的正式目标、课堂活动与核心评价按下表闭环。表中『答案证据』均位于本章练习的完整答案中,正文示例只承担示范作用。

表 12.1: 第十二章目标—活动—核心题映射
正式目标 正文活动 核心题与答案证据
表达式与优化 小节 12.2小节 12.4 练习 12.2 的表达式迁移
即时、延迟、流式与物化 小节 12.2.1小节 12.6 练习 12.1 与 12.4 的边界解释和一致性断言
Parquet 下推查询 小节 12.3小节 12.5 练习 12.3 的优化计划诊断
固定口径比较 Polars 与 Pandas 小节 12.7 练习 12.5 的结果一致性与比较协议审计

表 12.1 汇总本节的计算或审计结果,解释时应遵循正文给出的口径与限制。

12.2 列式执行与表达式代数

Polars 的 Python API 负责构造表达式,主要执行工作由 Rust 引擎完成。对一列 \(x=(x_1,\ldots,x_n)\),表达式 pl.col('close').log() 表示一个从输入列到输出列的变换,而不是逐行 Python 回调。若每行工作量近似恒定,向量化算子的工作量为 \(O(n)\);在 \(p\) 个有效线程上,理想计算时间下界可写为:

\[ T_p \geq \max\left(\frac{W}{p},\,T_{\mathrm{memory}},\,T_{\mathrm{serial}}\right), \tag{12.1}\]

式 12.1 给出本节后续实现与解释采用的数学关系。

其中 \(W\) 是总工作量,\(T_{\mathrm{memory}}\) 是内存带宽或 I/O 下界,\(T_{\mathrm{serial}}\) 是不能并行的部分。因此,多核实现并不意味着任意查询都会随核心数线性加速。

12.2.1 即时、延迟与流式的精确定义

表 12.2: Polars 三种执行概念的区别
模式 何时执行 优化范围 物化特点
即时 DataFrame 每次 API 调用时 单个算子附近 每一步返回已物化结果
延迟 LazyFrame collect() 等终端操作时 整个逻辑计划 执行器决定中间结果布局
流式执行引擎 终端操作触发后分批推进 支持流式的计划片段 降低部分中间结果峰值,不保证结果也超内存

表 12.2 汇总本节的计算或审计结果,解释时应遵循正文给出的口径与限制。

易混淆概念辨析:Lazy 不等于 Streaming

Lazy 指先构造逻辑计划、后优化执行;Streaming 指物理执行器尝试把计划分成批次推进。延迟查询可以由内存执行器执行,流式执行也必须先有可执行计划。全局排序、某些窗口、复杂连接或最终结果本身很大时,仍可能需要大量内存;是否支持还取决于具体 Polars 版本和算子。

12.3 本地真实 Parquet 数据与 schema

本章使用本地 full_market_with_return_2023.parquet 数据集。它由多个 Parquet 分片组成,字段是 order_book_idtrade_date、OHLC、volumetotal_turnoverdaily_return,而不是原稿中的模拟成交表。

import platform  # 识别Windows与Linux数据根目录
from pathlib import Path  # 构造可读的本地数据路径
from time import perf_counter  # 采用单调高精度时钟进行本机计时
import numpy as np  # 核验双引擎浮点聚合在数值容忍度内一致
import pandas as pd  # 提供同查询口径的基准实现
import polars as pl  # 构建表达式和延迟查询计划

DATA_ROOT = 'C:/qiufei/data' if platform.system() == 'Windows' else '/home/ubuntu/r2_data_mount/data'  # 选择规范数据根目录
MARKET_DATASET = Path(DATA_ROOT) / 'stock' / 'full_market_with_return_2023.parquet'  # 指向真实Parquet分片目录
MARKET_GLOB = str(MARKET_DATASET / '*.parquet')  # 让Polars扫描目录内全部分片
YRD_CODES = ['600104.XSHG', '600276.XSHG', '002415.XSHE', '002142.XSHE']  # 选择长三角代表公司
BENCHMARK_FILES = [str(MARKET_DATASET / file_name) for file_name in ['part.0.parquet', 'part.1.parquet', 'part.8.parquet', 'part.10.parquet']]  # 固定包含目标证券的四个真实物理分片
market_scan = pl.scan_parquet(MARKET_GLOB)  # 仅创建延迟扫描节点而不读取全部数据
market_schema = market_scan.collect_schema()  # 读取Parquet元数据以核对列名与类型
yrd_market_scan = market_scan.filter(pl.col('order_book_id').is_in(YRD_CODES) & (pl.col('trade_date') < pl.date(2024, 1, 1)))  # 固定四家公司2023年研究样本
compact_yrd_scan = pl.scan_parquet(BENCHMARK_FILES).filter(pl.col('order_book_id').is_in(YRD_CODES) & (pl.col('trade_date') < pl.date(2024, 1, 1)))  # 为重复练习限定包含样本的物理分片
print(market_schema)  # 展示真实数据schema
Schema([('order_book_id', String), ('trade_date', Datetime(time_unit='ns', time_zone=None)), ('open', Float64), ('high', Float64), ('low', Float64), ('adj_close', Float64), ('volume', Float64), ('total_turnover', Float64), ('daily_return', Float64), ('__null_dask_index__', Int64)])

读取 schema 通常只访问元数据;它不等同于执行全表查询。scan_parquet() 产生 LazyFrame,因此后续表达式可以参与全局优化。

12.4 表达式、筛选与聚合

下面的查询只保留四家公司,计算日收益、日内振幅,并汇总年度交易特征。表达式之间没有 Python 逐行循环,优化器可以把证券筛选和列裁剪推向 Parquet 扫描端。

yrd_query = (  # 构造尚未执行的年度统计计划
    yrd_market_scan  # 从已限定证券与年度的延迟扫描开始
    .select('order_book_id', 'trade_date', 'high', 'low', 'adj_close', 'volume', 'total_turnover', 'daily_return')  # 只保留分析需要的物理列
    .with_columns(((pl.col('high') - pl.col('low')) / pl.col('adj_close')).alias('intraday_range'))  # 定义相对日内振幅
    .group_by('order_book_id')  # 按证券汇总年度观测
    .agg(pl.len().alias('trading_days'), pl.col('daily_return').mean().alias('mean_daily_return'), pl.col('daily_return').std().alias('daily_volatility'), pl.col('intraday_range').mean().alias('mean_intraday_range'), pl.col('total_turnover').sum().alias('annual_turnover'))  # 一次声明多个可并行聚合
    .sort('annual_turnover', descending=True)  # 对小型聚合结果排序
)
print(yrd_query.explain(optimized=True))  # 输出优化后的物理相关计划证据
SORT BY [descending: [true]] [col("annual_turnover")]
  AGGREGATE[maintain_order: false]
    [len().alias("trading_days"), col("daily_return").mean().alias("mean_daily_return"), col("daily_return").std().alias("daily_volatility"), col("intraday_range").mean().alias("mean_intraday_range"), col("total_turnover").sum().alias("annual_turnover")] BY [col("order_book_id")]
    FROM
    simple π 4/4 ["order_book_id", ... 3 other columns]
       WITH_COLUMNS:
       [((col("high") - col("low")) / col("adj_close")).alias("intraday_range")] 
        simple π 6/6 ["order_book_id", "high", "low", ... 3 other columns]
          Parquet SCAN [/home/ubuntu/r2_data_mount/data/stock/full_market_with_return_2023.parquet/part.0.parquet, ... 15 other sources]
          PROJECT 7/10 COLUMNS
          SELECTION: (col("order_book_id").is_in([["600104.XSHG", "600276.XSHG", … "002142.XSHE"]]) & (col("trade_date") < 2024-01-01 00:00:00))
          ESTIMATED ROWS: 3644224

select() 是投影,filter() 是谓词,with_columns() 构造派生表达式,group_by().agg() 则把行映射为组统计量。优化器可改变安全算子的执行顺序,但必须保持结果语义。

yrd_result = yrd_query.collect()  # 使用默认执行引擎物化小型聚合结果
yrd_result  # 展示四家公司年度统计
表 12.3: Polars对长三角公司2023年真实行情的聚合结果
shape: (4, 6)
order_book_id trading_days mean_daily_return daily_volatility mean_intraday_range annual_turnover
str u32 f64 f64 f64 f64
"600276.XSHG" 241 0.000869 0.019035 0.026979 3.7818e11
"002415.XSHE" 241 0.000346 0.020718 0.030899 3.2866e11
"002142.XSHE" 241 -0.001709 0.017976 0.025389 1.9936e11
"600104.XSHG" 241 -0.000106 0.010892 0.015661 6.9623e10

12.4.1 结果解释

表 12.3 同时给出交易日数、日收益均值、日收益标准差、平均日内振幅和全年成交额。标准差与日内振幅衡量不同风险维度,不能互相替代;成交额较高只说明样本期交易更活跃,不构成公司质量或未来收益的因果判断。所有列都来自同一真实公司年度切片,因此还可以用交易日数识别停牌或数据覆盖差异。

12.4.2 窗口表达式与顺序约束

滚动收益依赖证券内部的时间顺序。over() 只声明分组范围,并不自动保证输入已按日期排序,因此应显式排序。

rolling_query = (  # 构建窗口分析计划
    yrd_market_scan  # 复用同一证券年度延迟扫描
    .select('order_book_id', 'trade_date', 'daily_return')  # 投影窗口所需三列
    .sort('order_book_id', 'trade_date')  # 显式建立证券内部时间顺序
    .with_columns(pl.col('daily_return').rolling_std(window_size=20).over('order_book_id').alias('volatility_20d'))  # 在证券分组内计算20期样本标准差
)
rolling_tail = rolling_query.filter(pl.col('volatility_20d').is_not_null()).tail(8).collect()  # 只物化有完整窗口的末尾少量结果
rolling_tail  # 展示窗口输出
shape: (8, 4)
order_book_id trade_date daily_return volatility_20d
str datetime[ns] f64 f64
"600276.XSHG" 2023-12-20 00:00:00 -0.012282 0.014446
"600276.XSHG" 2023-12-21 00:00:00 0.005878 0.014564
"600276.XSHG" 2023-12-22 00:00:00 -0.010789 0.014666
"600276.XSHG" 2023-12-25 00:00:00 -0.000227 0.01468
"600276.XSHG" 2023-12-26 00:00:00 -0.000455 0.014082
"600276.XSHG" 2023-12-27 00:00:00 0.002501 0.014054
"600276.XSHG" 2023-12-28 00:00:00 0.014969 0.01447
"600276.XSHG" 2023-12-29 00:00:00 0.010726 0.013357

12.5 查询优化与复杂度边界

\(n\) 行数据,筛选和投影通常为 \(O(n)\),哈希分组平均为 \(O(n)\)、内存约为 \(O(g)\),排序为 \(O(n\log n)\) 且常需较大中间缓冲区。Parquet 的统计元数据允许跳过不可能满足谓词的 row groups,但只有谓词与物理列兼容且统计信息可用时才能发生。

优化器常见动作包括:

  • 谓词下推:尽可能在扫描阶段排除不相关行;
  • 投影下推:只解码查询需要的列;
  • 表达式简化与公共子表达式复用:减少重复计算;
  • 切片下推:在语义允许时减少读取行数。

这些是计划级可能性,不应在没有查看 explain() 和执行指标时宣称一定发生。

12.6 Streaming 的能力与限制

筛选和分组聚合适合批次推进。为避免在同一教材内核中重复扫描全市场全部分片,下面用该真实数据集的一个物理分片验证两种引擎;engine='streaming' 请求流式引擎。它降低中间数据峰值的可能性,但若最终结果有数千万行,collect() 返回的 DataFrame 仍须放入内存。

streaming_example_query = pl.scan_parquet(str(MARKET_DATASET / 'part.0.parquet')).filter(pl.col('trade_date') < pl.date(2024, 1, 1)).group_by('order_book_id').agg(pl.col('total_turnover').sum().alias('annual_turnover'))  # 在一个真实分片上构建可流式聚合
streaming_baseline = streaming_example_query.collect()  # 使用默认引擎得到同分片基准
streaming_result = streaming_example_query.collect(engine='streaming')  # 请求流式物理执行同一查询
streaming_comparison = streaming_result.join(streaming_baseline, on='order_book_id', suffix='_baseline').with_columns(((pl.col('annual_turnover') - pl.col('annual_turnover_baseline')).abs() / pl.col('annual_turnover_baseline').abs()).alias('relative_difference'))  # 对齐两种引擎并量化浮点归约差异
assert streaming_comparison['relative_difference'].max() < 1e-12  # 在数值容忍度内验证两种物理引擎一致
print(streaming_result)  # 输出流式执行得到的小型最终结果
shape: (315, 2)
┌───────────────┬─────────────────┐
│ order_book_id ┆ annual_turnover │
│ ---           ┆ ---             │
│ str           ┆ f64             │
╞═══════════════╪═════════════════╡
│ 002658.XSHE   ┆ 1.7030e10       │
│ 300631.XSHE   ┆ 1.5497e10       │
│ 688711.XSHG   ┆ 2.9414e10       │
│ 603159.XSHG   ┆ 5.8950e9        │
│ 688613.XSHG   ┆ 1.2241e10       │
│ …             ┆ …               │
│ 002020.XSHE   ┆ 4.7220e10       │
│ 002218.XSHE   ┆ 2.2990e10       │
│ 002287.XSHE   ┆ 9.7125e9        │
│ 000721.XSHE   ┆ 1.7467e11       │
│ 688662.XSHG   ┆ 3.3392e10       │
└───────────────┴─────────────────┘

若目标是写出大结果,应优先考虑 sink_parquet() 等流式 sink,而不是先 collect() 再写盘。是否能完整流式执行必须结合当前版本的查询计划验证。

12.7 可复现的本机性能比较

性能不是库的固定常数。比较必须固定数据、列、筛选、聚合、缓存状态、线程数与输出物化方式。下面仅比较同一真实 Parquet 数据集上的一次端到端冷/暖状态混合测量;结果用于展示方法,不用于宣称 Polars 普遍快若干倍。

pandas_start_seconds = perf_counter()  # 记录Pandas端到端起始时刻
pandas_market = pd.read_parquet(BENCHMARK_FILES, columns=['order_book_id', 'daily_return', 'total_turnover'], filters=[[('order_book_id', 'in', YRD_CODES), ('trade_date', '<', pd.Timestamp('2024-01-01'))]])  # 在固定物理分片上请求相同证券、年度和列
pandas_result = pandas_market.groupby('order_book_id').agg(trading_days=('daily_return', 'size'), mean_daily_return=('daily_return', 'mean'), daily_volatility=('daily_return', 'std'), annual_turnover=('total_turnover', 'sum')).sort_values('annual_turnover', ascending=False)  # 执行同口径聚合
pandas_elapsed_seconds = perf_counter() - pandas_start_seconds  # 计算包含读取和物化的耗时
print(f'Pandas本次耗时为 {pandas_elapsed_seconds:.4f} 秒')  # 报告当前机器的一次实测
Pandas本次耗时为 9.0516 秒
polars_start_seconds = perf_counter()  # 记录Polars端到端起始时刻
polars_benchmark_query = pl.scan_parquet(BENCHMARK_FILES).filter(pl.col('order_book_id').is_in(YRD_CODES) & (pl.col('trade_date') < pl.date(2024, 1, 1))).group_by('order_book_id').agg(pl.len().alias('trading_days'), pl.col('daily_return').mean().alias('mean_daily_return'), pl.col('daily_return').std().alias('daily_volatility'), pl.col('total_turnover').sum().alias('annual_turnover')).sort('annual_turnover', descending=True)  # 在相同四个分片上声明同口径查询
polars_benchmark = polars_benchmark_query.collect()  # 执行与Pandas公共指标一致的查询
polars_elapsed_seconds = perf_counter() - polars_start_seconds  # 计算包含扫描和物化的耗时
benchmark_table = pd.DataFrame({'engine': ['Pandas', 'Polars'], 'elapsed_seconds': [pandas_elapsed_seconds, polars_elapsed_seconds], 'rows_returned': [len(pandas_result), polars_benchmark.height]})  # 汇总本次测量条件和结果
benchmark_table  # 输出可复核而不可过度外推的计时表
表 12.4: 当前环境下Pandas与Polars同口径查询的一次实测
engine elapsed_seconds rows_returned
0 Pandas 9.051576 4
1 Polars 43.229677 4

操作系统页缓存、首次加载动态库和 Polars 线程池初始化都会影响单次测量。严谨研究应交替顺序、多次重复、报告中位数与离散程度,并固定 POLARS_MAX_THREADS。本章删除原稿无依据的 5–50 倍断言。

12.8 Pandas 互操作与复制条件

Polars 可从 Pandas 或 Arrow 导入,也可导出 Pandas。是否共享缓冲区取决于源类型、目标后端、字符串表示、null 掩码、chunk 布局和可写性要求。普通 Pandas object 字符串通常要编码复制;使用 to_pandas(use_pyarrow_extension_array=True) 更有机会让兼容列共享 Arrow 缓冲区,但后续不支持 Arrow 的 Pandas 操作仍可能触发转换。

pandas_arrow_result = yrd_result.to_pandas(use_pyarrow_extension_array=True)  # 请求用Arrow扩展数组承接兼容列
print(pandas_arrow_result.dtypes)  # 检查实际输出后端而不假定零拷贝
print(pandas_arrow_result.head())  # 展示互操作后的业务结果
order_book_id          large_string[pyarrow]
trading_days                 uint32[pyarrow]
mean_daily_return            double[pyarrow]
daily_volatility             double[pyarrow]
mean_intraday_range          double[pyarrow]
annual_turnover              double[pyarrow]
dtype: object
  order_book_id  trading_days  mean_daily_return  daily_volatility  \
0   600276.XSHG           241           0.000869          0.019035   
1   002415.XSHE           241           0.000346          0.020718   
2   002142.XSHE           241          -0.001709          0.017976   
3   600104.XSHG           241          -0.000106          0.010892   

   mean_intraday_range      annual_turnover  
0             0.026979       378183021612.0  
1             0.030899  328658641639.950012  
2             0.025389   199361198650.48999  
3             0.015661        69623112680.0  

12.9 常见误区

  1. 把 Polars 描述为始终快于 Pandas。 小数据、已缓存数据或单个轻量操作可能由启动与转换开销主导。
  2. 把 Lazy 当成尚未占用任何内存。 逻辑计划本身很小,但执行时仍要读取、解码并物化必要状态。
  3. 认为 Streaming 可以执行任何超内存查询。 全局排序、复杂窗口和巨大最终结果可能成为阻塞点。
  4. 逐行使用 Python UDF。 这会阻止许多优化并引入 Python 调用开销,应优先组合原生表达式。
  5. 忽略行顺序。 Polars 没有 Pandas 隐式索引;窗口、滞后与累计计算前应明确排序键。

12.10 本章小结

Polars 的优势来自表达式、列式执行、多线程与查询优化的组合,而不是一个可脱离任务的性能倍数。Lazy 提供计划级优化机会,Streaming 是物理执行策略,两者必须分别定义。真实数据上的 explain()、结果一致性检查和可复现计时,才是选择工具的可靠证据。

12.11 分层练习与完整答案

12.11.1 练习 12.1:概念检查

说明为什么先 collect() 全表再筛选无法获得谓词下推,并指出 Streaming 仍可能内存不足的两个原因。

答案

collect() 已经执行扫描并物化全表,后续筛选只作用于内存中的 DataFrame,过滤条件无法回到 Parquet 扫描节点。Streaming 仍可能内存不足,因为计划包含全局排序等阻塞算子,或因为最终返回结果本身超过可用内存。

12.11.2 练习 12.2:表达式迁移

用 Polars Lazy API 计算四家公司正收益日比例,并按比例降序排列。

答案

positive_ratio = (  # 构建正收益日统计计划
    compact_yrd_scan  # 从固定物理分片的真实公司年度行情开始
    .group_by('order_book_id')  # 在证券层面汇总
    .agg((pl.col('daily_return') > 0).mean().alias('positive_day_ratio'))  # 布尔均值等于正收益日比例
    .sort('positive_day_ratio', descending=True)  # 将较高比例公司置于前部
    .collect()  # 执行并物化四行结果
)
positive_ratio  # 输出正收益日比例
shape: (4, 2)
order_book_id positive_day_ratio
str f64
"600276.XSHG" 0.473029
"600104.XSHG" 0.46888
"002415.XSHE" 0.435685
"002142.XSHE" 0.381743

12.11.3 练习 12.3:真实数据与计划诊断

计算每家公司成交额最高的五个交易日,并查看优化计划。解释该任务为何不能只用证券级聚合完成。

答案

top_turnover_query = (  # 构建公司内排序任务
    compact_yrd_scan  # 扫描固定物理分片的真实公司年度行情
    .select('order_book_id', 'trade_date', 'total_turnover')  # 下推列投影
    .sort(['order_book_id', 'total_turnover'], descending=[False, True])  # 建立公司内成交额顺序
    .group_by('order_book_id', maintain_order=True)  # 保持排序后的公司出现顺序
    .head(5)  # 每家公司保留前五条真实交易日记录
)
print(top_turnover_query.explain(optimized=True))  # 检查谓词和投影是否靠近扫描
top_turnover_query.collect()  # 输出二十条公司日记录
EXPLODE [trade_date, total_turnover]
  AGGREGATE[maintain_order: true]
    [col("trade_date").slice(offset=0, length=5).implode(), col("total_turnover").slice(offset=0, length=5).implode()] BY [col("order_book_id")]
    FROM
    SORT BY [descending: [false, true]] [col("order_book_id"), col("total_turnover")]
      Parquet SCAN [/home/ubuntu/r2_data_mount/data/stock/full_market_with_return_2023.parquet/part.0.parquet, ... 3 other sources]
      PROJECT 3/10 COLUMNS
      SELECTION: (col("order_book_id").is_in([["600104.XSHG", "600276.XSHG", … "002142.XSHE"]]) & (col("trade_date") < 2024-01-01 00:00:00))
      ESTIMATED ROWS: 911056
shape: (20, 3)
order_book_id trade_date total_turnover
str datetime[ns] f64
"002142.XSHE" 2023-10-12 00:00:00 3.0223e9
"002142.XSHE" 2023-07-25 00:00:00 2.5248e9
"002142.XSHE" 2023-07-28 00:00:00 2.1150e9
"002142.XSHE" 2023-07-31 00:00:00 2.0789e9
"002142.XSHE" 2023-04-24 00:00:00 2.0061e9
"600276.XSHG" 2023-07-31 00:00:00 6.2263e9
"600276.XSHG" 2023-01-16 00:00:00 5.8567e9
"600276.XSHG" 2023-08-01 00:00:00 5.8375e9
"600276.XSHG" 2023-08-07 00:00:00 5.2809e9
"600276.XSHG" 2023-01-19 00:00:00 4.7573e9

证券级求和会丢失交易日维度,而题目要求保留组内具体行,因此需要组内排序或 top-k 算子。

12.11.4 练习 12.4:综合挑战

分别用默认和流式引擎计算月度成交额,验证结果一致,并解释最终结果规模为何决定 collect() 是否安全。

答案

monthly_query = (  # 构建公司月度成交额计划
    compact_yrd_scan  # 复用固定物理分片的证券年度扫描
    .with_columns(pl.col('trade_date').dt.truncate('1mo').alias('trade_month'))  # 将交易日映射到自然月
    .group_by('order_book_id', 'trade_month')  # 形成证券月份分组
    .agg(pl.col('total_turnover').sum().alias('monthly_turnover'))  # 汇总月成交额
    .sort('order_book_id', 'trade_month')  # 建立稳定输出顺序
)
monthly_default = monthly_query.collect()  # 使用默认引擎得到基准结果
monthly_streaming = monthly_query.collect(engine='streaming')  # 使用流式引擎执行同一逻辑计划
monthly_comparison = monthly_default.join(monthly_streaming, on=['order_book_id', 'trade_month'], suffix='_streaming').with_columns(((pl.col('monthly_turnover') - pl.col('monthly_turnover_streaming')).abs() / pl.col('monthly_turnover').abs()).alias('relative_difference'))  # 按公司月比较两种浮点归约结果
assert monthly_comparison['relative_difference'].max() < 1e-12  # 在数值容忍度内验证两种物理策略一致
monthly_streaming.head(12)  # 展示一个公司的月度序列
shape: (12, 3)
order_book_id trade_month monthly_turnover
str datetime[ns] f64
"002142.XSHE" 2023-01-01 00:00:00 1.4007e10
"002142.XSHE" 2023-02-01 00:00:00 1.8189e10
"002142.XSHE" 2023-03-01 00:00:00 2.0907e10
"002142.XSHE" 2023-04-01 00:00:00 2.3063e10
"002142.XSHE" 2023-05-01 00:00:00 1.6373e10
"002142.XSHE" 2023-08-01 00:00:00 2.0313e10
"002142.XSHE" 2023-09-01 00:00:00 1.4627e10
"002142.XSHE" 2023-10-01 00:00:00 1.3124e10
"002142.XSHE" 2023-11-01 00:00:00 1.2438e10
"002142.XSHE" 2023-12-01 00:00:00 1.6295e10

四家公司一年仅约 48 个公司月,最终结果很小,collect() 安全;若按全市场逐笔明细返回,最终结果自身就可能超过内存,Streaming 也无法消除该物化成本。

12.11.5 练习 12.5:固定协议的双引擎比较

审计 表 12.4:先验证两种引擎在证券、交易日数和三个数值聚合上结果一致,再列出本次比较已经固定和仍未固定的条件。解释为什么这张表不能推出跨机器、跨查询的固定加速倍数。

答案

表 12.5 汇总本节的计算或审计结果,解释时应遵循正文给出的口径与限制。

pandas_comparable = pandas_result.sort_index()  # 将Pandas结果按证券键建立稳定比较顺序
polars_comparable = polars_benchmark.to_pandas().set_index('order_book_id').sort_index()  # 将小型Polars结果转为同索引表
assert pandas_comparable.index.astype(str).tolist() == polars_comparable.index.astype(str).tolist()  # 核验证券值相同并忽略Arrow桥接的字符串存储dtype
np.testing.assert_array_equal(pandas_comparable['trading_days'].to_numpy(), polars_comparable['trading_days'].to_numpy())  # 核验分组行数口径一致
comparison_columns = ['mean_daily_return', 'daily_volatility', 'annual_turnover']  # 固定需要容忍浮点归约误差的公共指标
np.testing.assert_allclose(pandas_comparable[comparison_columns], polars_comparable[comparison_columns], rtol=1e-12, atol=1e-12)  # 核验两个引擎的数值语义一致
protocol_audit = pd.DataFrame({'条件': ['物理分片', '证券与日期', '读取列与聚合', '输出物化', '线程数', '重复与顺序'], '状态': ['已固定', '已固定', '已固定', '已固定为小表', '未显式固定', '仅单次,未交替'], '含义': ['同一BENCHMARK_FILES', '同一四家公司和2023年上界', '同一四项聚合', '均返回内存结果', '只能报告当前环境', '不能估计计时离散度']})  # 披露可比性与剩余限制
protocol_audit  # 输出决定结论外推范围的协议清单
表 12.5: 练习12.5的双引擎结果与协议审计
条件 状态 含义
0 物理分片 已固定 同一BENCHMARK_FILES
1 证券与日期 已固定 同一四家公司和2023年上界
2 读取列与聚合 已固定 同一四项聚合
3 输出物化 已固定为小表 均返回内存结果
4 线程数 未显式固定 只能报告当前环境
5 重复与顺序 仅单次,未交替 不能估计计时离散度

固定输入、查询和物化方式只能使『本次任务』可比;线程池、页缓存、库版本、硬件、冷热顺序与重复次数仍会改变耗时。因此答案只能报告当前环境的一次观察,不能把耗时比写成 Polars 对任意 Pandas 工作负载的固定倍数。