13  DuckDB:SQL驱动的分析

13.1 先修知识与学习目标

本章假定读者掌握 SELECTWHEREGROUP BY 和基本连接,并理解 Parquet 的列式结构。DuckDB 是嵌入分析程序的进程内 OLAP 引擎,不是面向大量并发事务的服务器数据库,也不是把所有输入自动变成零拷贝的魔法层。

完成本章后,读者应能:

  • 目标 13.1:比较 OLTP 与 OLAP 的访问模式,并解释向量化执行为何适合扫描聚合;
  • 目标 13.2:直接查询本地真实 Parquet 分片,并观察投影、谓词与 row-group 跳过的条件;
  • 目标 13.3:将 HDF5 选择性读取结果显式转换为 Arrow 后注册给 DuckDB;
  • 目标 13.4:精确定义 DataFrame 互操作中的零拷贝成立条件、复制点与结果物化边界;
  • 目标 13.5:用窗口函数和持久化关系表达可复核的金融分析,并界定事务与并发边界。

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

表 13.1: 第十三章教学闭环映射
正式目标 学习活动或示例 可观察产出 核心练习 完整答案中的评分证据
目标 13.1 小节 13.2 的工作负载表与成本分解 为点查写入和扫描聚合场景选择引擎并说明机制 小节 13.9.5 OLTP/OLAP 场景判断、行式/列式与并发理由完整计分
目标 13.2 小节 13.3小节 13.3.1小节 13.3.3 输出 SQL 聚合结果及包含扫描、过滤、投影节点的计划 小节 13.9.2小节 13.9.3 负收益日 SQL、跨格式覆盖表与数值断言构成答案证据
目标 13.3 小节 13.4 的 HDF—Pandas—Arrow—DuckDB 链路 注册 Arrow 关系并核对 schema、行数和聚合 小节 13.9.3 共同证券—日期覆盖及成交额一致断言可自动评分
目标 13.4 小节 13.4.1 的五项成立条件 逐边标注共享、转换、算子状态和结果物化 小节 13.9.1小节 13.9.3 概念反例与 .df()、排序、字符串转换的复制点解释完整
目标 13.5 小节 13.5小节 13.6 输出并列极值窗口结果、持久表重开结果及回滚审计 小节 13.9.4小节 13.9.5 DENSE_RANK()、事务内计数、回滚后计数和重连计数均有断言

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

13.2 OLAP、向量化执行与资源边界

OLTP 负载通常以主键定位少量行并频繁更新;OLAP 负载则扫描较多行、读取少数列并执行聚合、排序和连接。DuckDB 按向量批次执行算子,而不是为每个单元格往返 Python。设扫描 \(n\) 行、选取 \(q\) 列、形成 \(g\) 个组,则简单哈希聚合的平均工作量约为 \(O(nq)\),状态内存约为 \(O(g)\);全局排序通常为 \(O(n\log n)\)

\[ T_{\mathrm{query}} \approx T_{\mathrm{scan}}+T_{\mathrm{decode}}+T_{\mathrm{operators}}+T_{\mathrm{materialize}}. \tag{13.1}\]

式 13.1 强调端到端时间还包含输出物化。即使扫描和聚合能够流式推进,巨大结果集转换成 Pandas 仍会占用与结果规模相称的内存。DuckDB 可在资源紧张时将部分排序或连接状态 spill 到磁盘,但磁盘空间、临时目录和具体算子仍构成边界。

表 13.2: OLTP 与嵌入式 OLAP 的工作负载差异
维度 典型 OLTP 系统 DuckDB 这类嵌入式 OLAP
主要任务 点查、插入、更新、小事务 扫描、连接、聚合、窗口
并发目标 多客户端并发事务 单进程分析吞吐
数据布局倾向 行式与索引定位 列式扫描与向量化执行
部署 长期运行的服务 应用进程内库或数据库文件

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

13.3 直接查询本地真实 Parquet

本章查询本地 2023 年全市场日行情 Parquet 分片。第一段代码建立跨平台路径并创建内存连接;不生成模拟 CSV,也不要求先将全部分片载入 Pandas。

import platform  # 识别操作系统以选择规范数据根目录
from pathlib import Path  # 以路径对象描述本地文件
import duckdb  # 提供进程内SQL分析引擎
import pandas as pd  # 选择性读取HDF5并核验结果
import pyarrow as pa  # 构造显式Arrow互操作边界

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')  # 构造DuckDB可扫描的分片模式
PRICE_HDF_PATH = Path(DATA_ROOT) / 'stock' / 'stock_price_pre_adjusted.h5'  # 指向真实前复权HDF5行情
YRD_CODES = ['600104.XSHG', '600276.XSHG', '002415.XSHE', '002142.XSHE']  # 固定四家长三角公司
connection = duckdb.connect(database=':memory:')  # 创建不写入仓库的内存数据库连接
print(f'DuckDB版本为 {duckdb.__version__}')  # 记录实际执行版本
DuckDB版本为 1.5.0

先通过 SQL 检查 schema。Parquet 元数据保留物理类型;字段名称表明该数据是行情,而非宏观指标。

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

# 查询 DuckDB 对真实 Parquet 扫描结果推断的字段名称与类型
parquet_schema = connection.sql(f'''
    -- 只描述扫描结果而不物化全市场明细
    DESCRIBE SELECT *
    -- 通配符使所有真实Parquet分片形成一个关系
    FROM read_parquet('{MARKET_GLOB}')
''').df()  # 将很小的schema结果转换为Pandas表
parquet_schema  # 展示字段名、SQL类型和可空性
表 13.3: DuckDB识别的本地真实行情Parquet schema
column_name column_type null key default extra
0 order_book_id VARCHAR YES None None None
1 trade_date TIMESTAMP_NS YES None None None
2 open DOUBLE YES None None None
3 high DOUBLE YES None None None
4 low DOUBLE YES None None None
5 adj_close DOUBLE YES None None None
6 volume DOUBLE YES None None None
7 total_turnover DOUBLE YES None None None
8 daily_return DOUBLE YES None None None
9 __null_dask_index__ BIGINT YES None None None

13.3.1 投影、谓词下推与聚合

下面的 SQL 仅选择四家公司,并只引用收益率、振幅与成交额所需列。Parquet 投影下推可减少列解码;谓词能否跳过 row group 取决于文件统计信息与物理排列,不能仅因写了 WHERE 就保证跳过全部无关字节。

# 汇总每只长三角证券的交易覆盖、日收益与成交额口径
yrd_aggregate = connection.sql(f'''
    -- 以证券为年度统计单元
    SELECT order_book_id,
           COUNT(*) AS trading_days,
           AVG(daily_return) AS mean_daily_return,
           STDDEV_SAMP(daily_return) AS daily_volatility,
           AVG((high - low) / NULLIF(adj_close, 0)) AS mean_intraday_range,
           SUM(total_turnover) AS annual_turnover
    -- 直接扫描真实Parquet分片
    FROM read_parquet('{MARKET_GLOB}')
    -- 只保留四家长三角上市公司
    WHERE order_book_id IN ('600104.XSHG', '600276.XSHG', '002415.XSHE', '002142.XSHE')
      AND trade_date < DATE '2024-01-01'
    -- 将可合并统计量压缩为四个证券组
    GROUP BY order_book_id
    -- 仅对四行聚合结果执行排序
    ORDER BY annual_turnover DESC
''').df()  # 物化小型聚合结果供教材展示
yrd_aggregate  # 输出真实行情统计
表 13.4: DuckDB直接查询Parquet得到的长三角公司年度统计
order_book_id trading_days mean_daily_return daily_volatility mean_intraday_range annual_turnover
0 600276.XSHG 241 0.000869 0.019035 0.026979 3.781830e+11
1 002415.XSHE 241 0.000346 0.020718 0.030899 3.286586e+11
2 002142.XSHE 241 -0.001709 0.017976 0.025389 1.993612e+11
3 600104.XSHG 241 -0.000106 0.010892 0.015661 6.962311e+10

13.3.2 结果解释

表 13.4 把四家公司的收益、波动、振幅与成交活跃度放在同一证券年度层级。daily_volatility 是日收益率的样本标准差,mean_intraday_range 则反映日内高低价区间相对收盘价的平均宽度;二者数值接近时也不表示经济含义相同。SQL 输出是描述性证据,不应被解释为流动性对收益或风险的因果效应。

13.3.3 执行计划证据

EXPLAIN 比口头宣称更可靠。不同版本的计划文本可能不同,但应能看到 Parquet 扫描、过滤、投影与聚合节点。

# 获取相同年度成交额聚合查询的执行计划,核验过滤和投影下推
query_plan = connection.sql(f'''
    -- 请求逻辑与物理计划而不返回业务明细
    EXPLAIN SELECT order_book_id, SUM(total_turnover) AS annual_turnover
    -- 扫描本地真实Parquet数据集
    FROM read_parquet('{MARKET_GLOB}')
    -- 证券谓词可能利用文件或row-group统计信息
    WHERE order_book_id IN ('600104.XSHG', '600276.XSHG')
      AND trade_date < DATE '2024-01-01'
    -- 在证券层面聚合成交额
    GROUP BY order_book_id
''').fetchone()[1]  # 取得当前DuckDB版本的计划文本
print(query_plan)  # 展示可核对的优化证据
┌───────────────────────────┐
│       HASH_GROUP_BY       │
│    ────────────────────   │
│         Groups: #0        │
│    Aggregates: sum(#1)    │
│                           │
│       ~142,891 rows       │
└─────────────┬─────────────┘
┌─────────────┴─────────────┐
│         PROJECTION        │
│    ────────────────────   │
│       order_book_id       │
│       total_turnover      │
│                           │
│       ~145,768 rows       │
└─────────────┬─────────────┘
┌─────────────┴─────────────┐
│           FILTER          │
│    ────────────────────   │
│ ((order_book_id = '600104 │
│ .XSHG') OR (order_book_id │
│     = '600276.XSHG'))     │
│                           │
│       ~145,768 rows       │
└─────────────┬─────────────┘
┌─────────────┴─────────────┐
│        READ_PARQUET       │
│    ────────────────────   │
│         Function:         │
│        READ_PARQUET       │
│                           │
│        Projections:       │
│       order_book_id       │
│       total_turnover      │
│                           │
│          Filters:         │
│ optional: order_book_id IN│
│   ('600104.XSHG', '600276 │
│          .XSHG')          │
│ trade_date<'2024-01-01 00 │
│   :00:00'::TIMESTAMP_NS   │
│                           │
│       ~728,844 rows       │
└───────────────────────────┘

13.4 从 HDF5 到 Arrow 再到 DuckDB

DuckDB 不原生扫描 Pandas HDFStore 表。合理链路是先利用 HDF5 的可查询索引选择性读取,再显式转换为 Arrow 表并注册。这一过程展示真实 HDF 转换,但整条链路不是零拷贝:HDF5 解码到 Pandas 已分配内存,object 股票代码转为 Arrow 字符串通常还要编码复制。

HDF_FILTERS = [f'order_book_id in {YRD_CODES}', 'date >= Timestamp(\'2023-01-01\')', 'date <= Timestamp(\'2023-12-31\')']  # 在HDF层限制证券与日期
hdf_prices = pd.read_hdf(PRICE_HDF_PATH, key='data', where=HDF_FILTERS, columns=['close', 'volume', 'total_turnover']).reset_index()  # 只载入目标行列
assert set(hdf_prices.columns) == {'order_book_id', 'date', 'close', 'volume', 'total_turnover'}  # 显式核验真实HDF schema
print(hdf_prices.shape)  # 报告选择性读取后的规模
print(hdf_prices.dtypes)  # 检查进入Arrow转换前的Pandas类型
(968, 5)
order_book_id             object
date              datetime64[ns]
close                    float64
volume                   float64
total_turnover           float64
dtype: object
hdf_arrow_table = pa.Table.from_pandas(hdf_prices, preserve_index=False)  # 建立Pandas到Arrow的显式格式边界
connection.register('yrd_hdf_arrow', hdf_arrow_table)  # 将Arrow表注册为只读SQL关系
# 对注册后的 Arrow 关系执行证券级覆盖与价格统计
arrow_query_result = connection.sql('''
    -- 对Arrow关系执行证券级验证统计
    SELECT order_book_id,
           COUNT(*) AS trading_days,
           AVG(close) AS mean_close,
           SUM(total_turnover) AS annual_turnover
    -- 直接引用已注册的Arrow关系
    FROM yrd_hdf_arrow
    -- 形成证券级结果
    GROUP BY order_book_id
    -- 采用稳定的证券顺序
    ORDER BY order_book_id
''').df()  # 将小型结果物化为Pandas DataFrame
arrow_query_result  # 展示HDF到Arrow再到SQL的完整真实链路
order_book_id trading_days mean_close annual_turnover
0 002142.XSHE 242 24.697573 2.002444e+11
1 002415.XSHE 242 33.563488 3.296943e+11
2 600104.XSHG 242 13.797560 6.987081e+10
3 600276.XSHG 242 44.232798 3.791746e+11

13.4.1 Zero-copy 的成立条件

一次互操作可称为零拷贝,至少需要同时满足:

  1. 源数据已经采用接收方可识别的内存格式,如兼容的 Arrow primitive/string buffers;
  2. 数据类型、null 表示、offset 宽度和 chunk 布局不需要重编码、强制转换或合并;
  3. 接收方可以只读扫描现有缓冲区,并保持源对象生命周期足够长;
  4. 所执行的算子不需要生成新列或重排数据;
  5. 输出消费者接受现有格式,而不是要求 NumPy 可写数组或 Python 对象。

因此,DuckDB 扫描兼容 Arrow 输入时可能复用输入缓冲区;聚合、排序、字符串转换和 .df() 输出通常会创建新的状态或数组。零拷贝是某条边上的性质,不是 DuckDB、Pandas 或 Polars 的永久属性。

arrow_output = connection.sql('SELECT order_book_id, close FROM yrd_hdf_arrow ORDER BY order_book_id, date LIMIT 8').fetch_arrow_table()  # 请求Arrow结果而非Python对象行
print(type(arrow_output))  # 确认输出容器是Arrow Table
print(arrow_output.schema)  # 检查实际输出类型与null表示
<class 'pyarrow.lib.Table'>
order_book_id: string
close: double

即使输出容器是 Arrow,也不能仅凭类型断言缓冲区与输入地址相同;这里的排序会重排数据,因而预期产生新输出缓冲区。

13.5 窗口函数:保留行粒度的分析

聚合会把多行压缩为一行;窗口函数则在保留原行的同时使用分区上下文。下例计算证券内部 20 日移动均值与当日成交额排名。窗口的排序键必须显式写入 OVER 子句。

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

# 在公司内按交易日和成交额两个顺序分别计算移动均价与排名
window_result = connection.sql(f'''
    -- 返回窗口计算后的末段记录
    SELECT order_book_id, trade_date, adj_close, total_turnover,
           AVG(adj_close) OVER (PARTITION BY order_book_id ORDER BY trade_date ROWS BETWEEN 19 PRECEDING AND CURRENT ROW) AS moving_average_20d,
           RANK() OVER (PARTITION BY order_book_id ORDER BY total_turnover DESC) AS turnover_rank
    -- 直接扫描真实Parquet行情
    FROM read_parquet('{MARKET_GLOB}')
    -- 以海康威视作为浙江上市公司案例
    WHERE order_book_id = '002415.XSHE' AND trade_date < DATE '2024-01-01'
    -- 教材输出只展示日期较晚的少量记录
    QUALIFY trade_date >= DATE '2023-12-20'
    -- 结果按交易日排序便于解释
    ORDER BY trade_date
''').df()  # 物化少量窗口结果
window_result  # 展示窗口统计
表 13.5: DuckDB对真实行情计算20日均价与公司内成交额排名
order_book_id trade_date adj_close total_turnover moving_average_20d turnover_rank
0 002415.XSHE 2023-12-20 733.5786 5.171172e+08 763.950510 220
1 002415.XSHE 2023-12-21 735.3678 5.862563e+08 760.629280 205
2 002415.XSHE 2023-12-22 731.3420 6.330535e+08 757.151490 192
3 002415.XSHE 2023-12-25 736.7097 5.367252e+08 754.534760 217
4 002415.XSHE 2023-12-26 733.5786 5.356998e+08 751.761475 218
5 002415.XSHE 2023-12-27 736.7097 5.503862e+08 749.435495 214
6 002415.XSHE 2023-12-28 757.0620 9.937313e+08 748.384330 133
7 002415.XSHE 2023-12-29 776.5197 1.141917e+09 747.512085 113

ROWS BETWEEN 19 PRECEDING AND CURRENT ROW 是包含当前行的最多 20 行窗口;样本最前端不足 20 行时仍计算现有行均值。若业务要求必须满 20 个交易日,应同时计算窗口计数并过滤。

13.6 持久化数据库与事务边界

duckdb.connect('analysis.duckdb') 会创建持久化数据库文件,适合反复使用的清洗表、视图和统计结果;:memory: 连接适合一次性分析。持久化并不等同于多用户 OLTP 服务:同一数据库文件的并发写入模式和进程协调需要遵循 DuckDB 的并发约束。教材渲染使用内存连接,避免向仓库写入运行产物。

持久化回答“连接关闭后对象是否仍存在”,事务回答“一组变更是整体提交还是整体撤销”,两者不是同一维度。BEGIN TRANSACTION 后,本连接可以看到尚未提交的写入;ROLLBACK 必须撤销该事务的全部变更;COMMIT 才把变更提交到数据库状态。事务原子性也不把嵌入式数据库变成高并发支付服务:客户端并发、长事务冲突、服务可用性、访问控制与灾备仍需独立设计。小节 13.9.5 在临时目录创建文件数据库,既验证重连持久性与回滚,又不把运行产物写入仓库。

# 创建只保存查询定义的年度市场视图,供持久化与事务练习复用
connection.sql(f'''
    -- 视图保存查询定义而非复制全部行情
    CREATE OR REPLACE VIEW yrd_market_2023 AS
    -- 仅暴露后续分析需要的列
    SELECT order_book_id, trade_date, adj_close, daily_return, total_turnover
    -- 视图底层仍指向本地Parquet分片
    FROM read_parquet('{MARKET_GLOB}')
    -- 固定长三角公司研究样本
    WHERE order_book_id IN ('600104.XSHG', '600276.XSHG', '002415.XSHE', '002142.XSHE')
      AND trade_date < DATE '2024-01-01'
''')  # 在当前内存连接中注册可复用查询定义
view_row_count = connection.sql('SELECT COUNT(*) FROM yrd_market_2023').fetchone()[0]  # 验证视图能够访问真实数据
print(f'视图包含 {view_row_count} 条公司日记录')  # 报告视图规模
视图包含 964 条公司日记录

13.7 常见误区

  1. 把嵌入式等同于玩具数据库。 部署形态与分析能力是两个维度;DuckDB 能执行复杂 OLAP,但不以高并发事务为目标。
  2. 把 SQL 写法等同于执行方式。 是否下推、并行或 spill 必须查看计划和实际指标。
  3. 宣称 Pandas、Polars 与 DuckDB 无条件零拷贝。 object 字符串、类型转换、重排和结果物化都会触发复制。
  4. 先读入 Pandas 全表再交给 DuckDB。 这会丢失 Parquet 扫描阶段的列裁剪和 row-group 跳过机会。
  5. 对窗口省略排序键。 没有显式业务顺序的移动计算在金融上没有可解释性。

13.8 本章小结

DuckDB 把 SQL、列式文件扫描和向量化执行嵌入 Python 进程。其优势来自直接在数据源附近完成筛选、投影、聚合与窗口计算。HDF5 需要显式转换边界,Arrow 互操作只有在格式兼容、无需重排且消费者接受现有缓冲区时才可能零拷贝;任何定量结论都应由执行计划和实际物化路径支持。

13.9 分层练习与完整答案

13.9.1 练习 13.1:概念检查

判断下列说法是否正确并解释:DuckDB 查询一个 Pandas DataFrame,所以输入和输出都一定零拷贝。

答案

错误。DuckDB 可能直接扫描某些兼容输入缓冲区,但 object 字符串常需转换;过滤、聚合、排序会创建状态或新缓冲区;.df() 还要构造 Pandas 输出。必须分别审计输入扫描、算子和输出三段链路。

13.9.2 练习 13.2:SQL 聚合

用 SQL 计算四家公司负收益日数量与占比,按负收益日占比降序排列。

答案

# 汇总各证券有效收益日、下跌日数量及下跌日占有效收益日的比例
negative_day_result = connection.sql('''
    -- 按证券汇总下跌日风险频率
    SELECT order_book_id,
           COUNT(daily_return) AS valid_return_days,
           COUNT(*) FILTER (WHERE daily_return < 0) AS negative_days,
           COUNT(*) FILTER (WHERE daily_return < 0) * 1.0
               / NULLIF(COUNT(daily_return), 0) AS negative_day_ratio
    -- 使用已定义的长三角行情视图
    FROM yrd_market_2023
    -- 形成证券级统计
    GROUP BY order_book_id
    -- 将下跌频率较高的公司置于前部
    ORDER BY negative_day_ratio DESC
''').df()  # 物化四行练习结果;缺失收益不进入比例分母

assert (negative_day_result['negative_days'] <= negative_day_result['valid_return_days']).all()  # 下跌日必须是有效收益日的子集
negative_day_valid_rows = negative_day_result['valid_return_days'] > 0  # 只在存在有效收益时核验可定义的比例
pd.testing.assert_series_equal(                            # 直接核验 SQL 比例确实使用有效收益日而非总行数
    negative_day_result.loc[negative_day_valid_rows, 'negative_day_ratio'].reset_index(drop=True),
    (negative_day_result.loc[negative_day_valid_rows, 'negative_days']
     / negative_day_result.loc[negative_day_valid_rows, 'valid_return_days']).reset_index(drop=True),
    check_names=False,
)
assert negative_day_result.loc[negative_day_valid_rows, 'negative_day_ratio'].between(0, 1).all()  # 核验频率范围
negative_day_result  # 输出有效收益日、下跌日数量与比例
order_book_id valid_return_days negative_days negative_day_ratio
0 002142.XSHE 241 149 0.618257
1 002415.XSHE 241 134 0.556017
2 600276.XSHG 241 125 0.518672
3 600104.XSHG 241 123 0.510373

13.9.3 练习 13.3:Parquet 与 HDF 结果核验

审计 Parquet 视图和 HDF-Arrow 关系的公司日覆盖范围,并在共同的证券—日期观测上验证成交额一致。不得假设两个本地快照具有完全相同的起止日。

答案

# 对齐 Parquet 与 HDF 来源的公司记录,比较两种格式的年度成交额
format_comparison = connection.sql('''
    WITH parquet_rows AS (SELECT order_book_id, trade_date AS date, total_turnover AS parquet_turnover FROM yrd_market_2023),
         hdf_rows AS (SELECT order_book_id, date, total_turnover AS hdf_turnover FROM yrd_hdf_arrow),
         aligned_rows AS (SELECT COALESCE(parquet_rows.order_book_id, hdf_rows.order_book_id) AS order_book_id, parquet_turnover, hdf_turnover FROM parquet_rows FULL OUTER JOIN hdf_rows USING (order_book_id, date))
    SELECT order_book_id,
           COUNT(parquet_turnover) AS parquet_rows,
           COUNT(hdf_turnover) AS hdf_rows,
           COUNT(*) FILTER (WHERE parquet_turnover IS NOT NULL AND hdf_turnover IS NOT NULL) AS common_rows,
           MAX(ABS(parquet_turnover - hdf_turnover) / NULLIF(ABS(hdf_turnover), 0)) AS maximum_relative_difference
    FROM aligned_rows
    GROUP BY order_book_id
    ORDER BY order_book_id
''').df()  # 以全外连接同时审计覆盖差异和共同观测数值
assert format_comparison['maximum_relative_difference'].max() < 1e-10  # 仅在共同公司日上验证数值一致
assert format_comparison['common_rows'].min() > 0  # 确认每家公司都有可比较的共同观测
format_comparison  # 输出逐证券核验表
order_book_id parquet_rows hdf_rows common_rows maximum_relative_difference
0 002142.XSHE 241 242 241 0.0
1 002415.XSHE 241 242 241 0.0
2 600104.XSHG 241 242 241 0.0
3 600276.XSHG 241 242 241 0.0

输出会显示两个本地快照的行数可能略有差异;这正是跨格式或跨快照分析必须先审计键覆盖、再比较共同观测的原因。

13.9.4 练习 13.4:综合窗口任务

计算每家公司 2023 年最大单日跌幅及发生日期,并在跌幅并列时保留所有日期。

答案

# 使用窗口稠密排名提取每家公司样本内并列最大跌幅日
worst_days = connection.sql('''
    -- 先为每家公司按收益率从低到高建立稠密排名
    WITH ranked_returns AS (
        SELECT order_book_id, trade_date, daily_return,
               DENSE_RANK() OVER (PARTITION BY order_book_id ORDER BY daily_return ASC) AS loss_rank
        FROM yrd_market_2023
        WHERE daily_return IS NOT NULL
    )
    -- 保留排名第一的全部并列最差日期
    SELECT order_book_id, trade_date, daily_return
    FROM ranked_returns
    WHERE loss_rank = 1
    ORDER BY order_book_id, trade_date
''').df()  # 物化规模很小的极端收益结果
worst_days  # 输出最差交易日及收益率
order_book_id trade_date daily_return
0 002142.XSHE 2023-10-12 -0.055659
1 002415.XSHE 2023-06-09 -0.063386
2 600104.XSHG 2023-03-10 -0.044876
3 600276.XSHG 2023-07-31 -0.091131

使用 DENSE_RANK() 而不是 ROW_NUMBER(),可以在最小收益率并列时保留全部观测,符合题意。

13.9.5 练习 13.5:OLTP/OLAP 选型与持久化事务

一家券商同时有两个任务:A 为数万名客户并发提交委托并更新账户余额;B 为研究员扫描一亿行历史行情,按证券聚合风险指标。回答:

  1. 哪个任务更符合 OLTP,哪个更符合 DuckDB 的嵌入式 OLAP 定位?至少从访问行数、写入模式、并发目标和数据布局四个维度说明;
  2. 在临时文件数据库创建空的 audit_events 持久表,事务内插入一行后回滚;验证事务内计数为 1、回滚后为 0,并在关闭和重新连接后验证表仍存在且行数仍为 0;
  3. 解释为什么上述结果证明了“持久化对象 + 事务回滚”,却没有证明 DuckDB 适合任务 A。

评分点(100 分):OLTP/OLAP 判断 30 分;四维理由 20 分;事务内与回滚后断言 25 分;重连持久化断言 15 分;并发边界解释 10 分。

答案

任务 A 是典型 OLTP:每次按账户或订单键定位少量行,频繁短事务写入,并以大量客户端并发、低延迟和一致性为核心。任务 B 是典型 OLAP:扫描很多行、读取少数列、以聚合和排序为主,列式布局与向量化批处理能减少解释器往返并提高扫描吞吐。DuckDB 适合嵌入研究进程执行 B;A 通常需要面向并发事务的服务数据库及完整运维体系。

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

import tempfile  # 创建自动清理的隔离目录,避免教材渲染留下数据库文件
with tempfile.TemporaryDirectory(prefix='ch13_duckdb_') as exercise_directory:  # 将持久化实验限定在可自动回收的路径
    exercise_database_path = Path(exercise_directory) / 'audit.duckdb'  # 为关闭后重连准备真实文件数据库路径
    persistent_connection = duckdb.connect(str(exercise_database_path))  # 打开会跨连接保存schema的文件数据库
    persistent_connection.sql('CREATE TABLE audit_events (event_id INTEGER, event_name VARCHAR)')  # 在事务外提交空的持久表定义
    persistent_connection.sql('BEGIN TRANSACTION')  # 开始需要整体提交或撤销的一组写入
    persistent_connection.sql("INSERT INTO audit_events VALUES (1, 'price_check')")  # 插入一条仅在当前事务可见的审计事件
    rows_inside_transaction = persistent_connection.sql('SELECT COUNT(*) FROM audit_events').fetchone()[0]  # 核验事务内能看到待定写入
    persistent_connection.sql('ROLLBACK')  # 撤销本事务的插入但保留此前已提交的表定义
    rows_after_rollback = persistent_connection.sql('SELECT COUNT(*) FROM audit_events').fetchone()[0]  # 核验回滚后数据行已消失
    persistent_connection.close()  # 关闭首个连接以测试对象是否持久化
    reopened_connection = duckdb.connect(str(exercise_database_path))  # 从同一文件建立全新连接
    rows_after_reopen = reopened_connection.sql('SELECT COUNT(*) FROM audit_events').fetchone()[0]  # 核验重连后表存在且回滚状态被保存
    reopened_connection.close()  # 在临时目录清理前释放数据库文件句柄
assert (rows_inside_transaction, rows_after_rollback, rows_after_reopen) == (1, 0, 0)  # 固定事务可见性、回滚与重连的评分合同
pd.DataFrame({'checkpoint': ['事务内', '回滚后', '关闭并重连后'], 'row_count': [rows_inside_transaction, rows_after_rollback, rows_after_reopen]})  # 输出三阶段可观察证据
表 13.6: 练习13.5的持久化与事务回滚审计
checkpoint row_count
0 事务内 1
1 回滚后 0
2 关闭并重连后 0

表定义在事务开始前创建,因此关闭并重连后仍可查询;插入发生在显式事务内并被回滚,所以后两次计数均为 0。这个实验验证单文件中的持久化和原子撤销,不包含多客户端压力、故障转移、账户级锁竞争或服务可用性证据,不能据此把任务 A 交给嵌入式 OLAP 引擎。

connection.close()  # 释放DuckDB连接与已注册关系