下面的 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 # 输出真实行情统计
# 获取相同年度成交额聚合查询的执行计划,核验过滤和投影下推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) # 展示可核对的优化证据
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 Tableprint(arrow_output.schema) # 检查实际输出类型与null表示
# 在公司内按交易日和成交额两个顺序分别计算移动均价与排名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 个交易日,应同时计算窗口计数并过滤。
# 对齐 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 # 输出逐证券核验表
# 使用窗口稠密排名提取每家公司样本内并列最大跌幅日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 # 输出最差交易日及收益率
任务 A 是典型 OLTP:每次按账户或订单键定位少量行,频繁短事务写入,并以大量客户端并发、低延迟和一致性为核心。任务 B 是典型 OLAP:扫描很多行、读取少数列、以聚合和排序为主,列式布局与向量化批处理能减少解释器往返并提高扫描吞吐。DuckDB 适合嵌入研究进程执行 B;A 通常需要面向并发事务的服务数据库及完整运维体系。