Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
27f70f2
fix: ddb_dataset_processor 不再丢弃 inst_processors 处理结果
hugo2046 Jul 21, 2026
d088f8c
fix: DDBClient 连接池生命周期修复并移除死代码
hugo2046 Jul 21, 2026
41f9458
fix: DolphinDBClientProvider 延迟创建连接池
hugo2046 Jul 21, 2026
482654d
fix: load_mysql_plugin 安装 mysql 插件而非 lgbm
hugo2046 Jul 21, 2026
d3b44b9
fix: DolphinDBDataLoader 不再关闭进程级共享 session
hugo2046 Jul 21, 2026
84aea0d
fix: DDB 会话访问统一加会话级 RLock
hugo2046 Jul 21, 2026
b435008
perf: QlibDataLoader 全局 load 锁收窄至 DolphinDB 后端
hugo2046 Jul 21, 2026
1f3d85d
chore: DDB 后端日志与异常处理统一
hugo2046 Jul 21, 2026
9aab134
fix: MySQL 同步 SQL 参数校验(防注入)
hugo2046 Jul 21, 2026
bcb0849
chore: 清理 DDB 后端死代码与重复算子映射
hugo2046 Jul 21, 2026
8deb733
test: fetch_features_from_ddb 分支回归与 RPC 基线(D0 安全网)
hugo2046 Jul 21, 2026
615a26b
perf: 交易日历缓存(消除每批字段的全量日历下载)
hugo2046 Jul 21, 2026
f409dc0
perf: fetch_features_from_ddb 往返缩减与结构化拆分
hugo2046 Jul 21, 2026
c7a9820
perf: 存储层/schema/表达式翻译缓存
hugo2046 Jul 21, 2026
aae3881
perf: 计算分支结果直构 MultiIndex 面板
hugo2046 Jul 21, 2026
8358416
perf: alpha 因子库按需惰性加载
hugo2046 Jul 21, 2026
4c9c489
perf: 写路径按 URI 复用共享 DDBClient
hugo2046 Jul 21, 2026
7fa6f86
perf: DDB 批次参数配置化(默认值不变)
hugo2046 Jul 21, 2026
3a9cac0
docs: DDB 后端优化的 CHANGELOG 与缓存/调参说明
hugo2046 Jul 21, 2026
ab525e9
fix: DDB 计算分支自动取前序期并落实真实分段计算
hugo2046 Jul 22, 2026
6ee1096
fix: 修复分段合并丢弃正确分段与回看正则兜底误判常量
hugo2046 Jul 22, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
108 changes: 108 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,113 @@
# CHANGELOG

## 2026-07-22(审查修复)

### 修复代码审查发现的两个问题(分段丢弃 + 正则兜底误判)

- **分段合并丢弃正确分段**:`mergeSegmentResults` 旧实现遇到占位 TABLE(某
分段计算失败/空的兜底)就 `break`,把该表达式其余分段已算出的正确矩阵
全部丢弃,与「分段边界结果与整段一致」的保证矛盾。改为把失败/空分段的
占位由 `buildPlaceholderMatrix` 统一产出「行-该段交易日、列-presentCodes」
的 NULL 矩阵(与正常结果同为带标签矩阵):合并时一视同仁按行拼接,任一
分段异常都不再丢弃其它分段;顺带消除空 `baseData` 分段的日期空洞,并去掉
Python 侧期望三元素矩阵却收到 TABLE 的类型隐患。占位填 NULL 而非 0,避免
污染缺失位置的因子值。移除因此不再被调用的 `createEmptyTable`。
- **正则兜底把数值常量误判为窗口**:回看解析兜底(仅对 qlib 解析不了的表达式
生效)取最大独立整数为窗口,会把 `gtjaAlpha191_005($volume/1000000, ...)`
里的 `1000000` 当成百万日回看,在 2 核/8GB 上触发超量查询。新增
`_MAX_FALLBACK_WINDOW=2000`(约 8 年交易日)上界:超过即视为常量剔除,
剔除后无候选则退回 `ddb_lookback_default`。新增 3 项针对性单测。

## 2026-07-22

### DDB 计算分支:滚动算子自动取前序期 + 真实分段计算(替代虚假 repartitionDS)

修复两个相互纠缠的老问题(详见
`docs/DDB取前序期与分段计算_20260722.md`):

- **取前序期**:DDB 后端此前把用户的 `start_time` 原样下发,
`Mean($close,20)` 在区间头部前 19 个交易日必然 NaN(文件后端靠
`get_extended_window_size` 外扩,DDB 分支绕过了该机制)。现在 Python 侧
复用 qlib 算子树解析每批表达式的(向前, 向后)外扩交易日数(嵌套算子
递归累加、双臂算子取 max、`Ref($close,-2)` 标签向后外扩),经
`FeatureEngineeringByDate` 尾参传给 DDB;服务器端外扩查询窗口计算后
把结果矩阵截断回请求区间。qlib 无法实例化的表达式(DDB 专属 alpha 库
函数)退回启发式:正则扫独立整数当窗口,扫不到用
`C["ddb_lookback_default"]`(默认 252)兜底。
- **分段计算**:原 `repartitionDS(..., RANGE, [start, end])` 只有 2 个边界
→ 永远只产生 1 个数据源,mr 退化为单任务(代码中 FIXME 自证的
"虚假的分区")。且真把日期切段后每段头部都会缺窗口——与取前序期
必须一起设计。现移除 repartitionDS+mr,改为:内存估算
(`FeatureEngine.isRunWithAvailableMemory`,扩展窗口 + 表达式矩阵项 +
70% 空闲内存上限)放得下则整段单次计算;放不下按
`C["ddb_days_step"]` 切段**顺序循环**(单机 2 核 mr 并行收益小且峰值
内存翻倍),每段查询带 lookback 重叠(halo)、段内截断,各段 panel 用
统一列标签(窗口内实际出现的代码集合)对齐后 `concatMatrix` 纵向合并,
保证分段边界处滚动结果与整段计算一致。

接口兼容:`FeatureEngineeringByDate` 新增尾参
`lookbackDays=0, rightDays=0`(旧调用形式行为不变);Python 公共接口
`D.features()` 签名不变。离线单测 +12(回看解析/批量取 max/脚本接线),
DDB 端脚本行为待 live 验证。

## 2026-07-21

### DDB 后端系统性优化(分支 optimize/ddb-backend,17 个原子提交)

面向 DolphinDB **社区版(2 核 / 8GB)** 的一轮优化:省往返、省传输、
省服务器内存、控会话数;公共接口 `D.features()/D.calendar()/D.instruments()`
的签名与返回结果完全不变。全部改动配套离线单测(mock 会话,共 115 项通过)。

**Bug 修复**:

- `ddb_dataset_processor` 在传入 `inst_processors` 时提前 `return pd.DataFrame()`,
并行处理结果被整体丢弃(`data.py:782`)。
- `DDBClient` 连接池:`_pool_instance` 类变量导致多客户端共享同一池;
`close_pool` 引用不存在的 `_pool_lock` 且永远读到 None(从未真正关闭过池);
删除调用不存在 `get_session()` 的死方法 `tableAppender/tableUpsert` 与
`__main__` 中的硬编码凭据。
- `DolphinDBClientProvider` init 时急切创建 4 连接的池(无消费方)→ 改懒创建。
- `load_mysql_plugin` 误装 `lgbm` 插件(应为 `mysql`)。
- `DolphinDBDataLoader.__exit__/__del__` 无条件关闭进程级共享 session →
引入 `_owns_session` 所有权标志。

**并发/线程安全**:

- 新增会话级 `DBClient.session_lock`(RLock):feature 查询是
「run→upload→run」多步会话对话,交错执行会互相覆盖服务器变量
(跨线程数据污染根因);所有 session 触点统一持锁。
- `QlibDataLoader` 全局 `_load_lock` 收窄至 DDB 路径,文件后端恢复无锁并行。

**健壮性/质量**:

- 12 处 `print` → loguru;3 处裸 `except:` 收敛;异常重包装补 `from e`。
- MySQL 同步 SQL 参数白名单校验(`validate_date_str`/`validate_sql_identifier`,
唯一真实注入面)。
- 清理死代码(注释块 ×2、零调用者函数 ×2、16KB 备份脚本)与
`OPERATOR_MAPPING` 重复键(生效映射不变,ast 测试锁定)。

**性能**(RPC 往返:计算分支每批 3-4 次 → 预热后 2 次):

- 交易日历模块级缓存:Alpha158 一次 `D.features` 从 ~6 次全量日历下载 → 0;
`DBCalendarStorage.index()/__getitem__` 改走 `H["c"]` 缓存。
- 日期字面量内联(消灭独立上传往返);纯字段分支合并为单条 SQL 脚本;
上传字典按分支裁剪。`fetch_features_from_ddb` 拆为编排器 + 4 个 helper。
- existsTable 仅正向缓存、股票池走 `H["i"]`、表 schema 进程内缓存、
表达式翻译 lru_cache;统一由 `ddb_qlib.invalidate_ddb_caches()` 失效
(写路径自动调用)。
- 计算分支结果直构 `(instrument, datetime)` MultiIndex 面板,替代
concat→unstack→stack→swaplevel 的 ~4 次全景拷贝(保留 legacy 兜底 +
逐值等价测试)。
- 三个 alpha 因子库(约 119KB)按字段前缀惰性加载(`Cannot recognize the
token` 兜底重试;`C["ddb_preload_alpha_libs"]` 逃生开关)。
- 写路径按 URI 复用共享 `DDBClient`(批量导入 N 会话 → 1)。
- 批次参数配置化:`C["ddb_field_chunk_size"]=30`、`C["ddb_days_step"]=252`
(默认值不变,`scripts/benchmark_ddb_backend.py` 供 live 校准)。

**明确不做**(2 核/8GB 约束):读路径不引入 `DBConnectionPool` 并发、
不调高 mr 并行度、不做服务器端全景 pivot、不引入 module/functionView 持久化。


## 2026-07-20

### fix: DDB 批量取数路径恢复成分股 spans(入池/出池)过滤 ⚠️ BREAKING 行为变更
Expand Down
118 changes: 118 additions & 0 deletions docs/DDB取前序期与分段计算_20260722.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,118 @@
# DDB 后端:滚动算子取前序期与真实分段计算

> 分支 `optimize/ddb-backend`,2026-07-22。
> 涉及文件:`qlib/data/backend/ddb_qlib/ddb_features.py`、
> `qlib/data/backend/ddb_qlib/ddb_scripts/featureEngineering.dos`、`qlib/config.py`。

## 背景:两个纠缠在一起的老问题

### 问题一:DDB 后端没有取前序期

qlib 文件后端在 `LocalExpressionProvider.file_expression` 中先调用算子树的
`get_extended_window_size()` 得到左右外扩量,把查询窗口外扩后计算、再截断回
请求区间(`qlib/data/data.py`)。而 DDB 分支 `ddb_expression` 直接把用户的
`start_time`/`end_time` 下发给服务器端 `parseExpr` 计算——`Mean($close,20)`
在 `start_date` 后的前 19 个交易日必然是 NaN。

### 问题二:repartitionDS 分区从未生效

`FeatureEngine.createDSByDate` 写的是:

```dolphinscript
repartitionDS(sql(...), `date, RANGE, [startTime, boundEndDate])
```

RANGE 分区方案 n+1 个边界产生 n 个数据源——这里只有 2 个边界,**永远只产生
1 个数据源**,`mr` 退化为单任务:既没有并行,也没有分批控内存(代码中
FIXME "虚假的分区" 自证)。

### 为什么必须一起设计

一旦让日期分段真的生效,每个分段独立计算时滚动算子在**每一段的开头**都会
产生 L 天 NaN(不止第一段)。而 repartitionDS 的 RANGE 分区天然不重叠,
表达不了"每段多带 L 天头部数据",所以 repartitionDS+mr 架构对滚动算子
结构性不适用。

## 方案

### Python 侧:回看窗口解析(`ddb_features.py`)

`get_expression_extended_window(expr, default_lookback) -> (lft, rght)`:

1. **优先复用 qlib 算子树**:用 `parse_field` + 受限命名空间 eval 实例化
表达式(与 `ExpressionProvider.get_expression_instance` 同款),调
`get_extended_window_size()`。嵌套算子递归累加(`Mean(Ref($close,5),20)`
→ 24)、双臂算子取 max(`Corr`)、未来引用向后外扩(标签
`Ref($close,-2)` → rght=2)全部免费获得。
2. **启发式兜底**:qlib 无法实例化的表达式(gtjaAlpha/WQAlpha/qlib158Alpha
等 DDB 专属函数)退回正则——独立整数字面量当窗口(标识符里的数字如
`gtjaAlpha191_001` 不会误判);扫不到时用 `C["ddb_lookback_default"]`
(默认 252)兜底。真正非法的表达式仍在 DDB 端 parseExpr 阶段报错。

`batch_extended_window`:同批表达式共享一份基础数据面板,按批取 max。
结果经 `FeatureEngineeringByDate(..., daysStep, lookbackDays, rightDays)`
尾参下发(旧调用形式默认 0,行为不变)。

### DDB 侧:外扩 + 截断 + 真实分段(`featureEngineering.dos`)

`FeatureEngine.fetch(daysStep)` 流程:

1. `shiftTradeDays` 按 `temporalAdd(dt, ±N, `XSHG)` 计算扩展窗口
`[extStart, extEnd]`。
2. `fetchPresentCodes`:扩展窗口内实际有数据的代码集合(升序),作为
所有 panel 的统一列标签。
3. `isRunWithAvailableMemory(extStart, extEnd)`:峰值内存估算 ≈ 基础长表
×2(长表 + panel 字典)+ 每表达式一个 dates×codes 的 DOUBLE 矩阵,
上限为 70% 空闲内存。
4. **放得下** → `computeRangeFeatures(startTime, endTime, presentCodes)`
整段单次计算:查询 `[extStart, extEnd]`,panel + parseExpr 计算后用
`loc` 把结果矩阵行截断回 `[startTime, endTime]`(外扩部分只服务于
窗口计算,不进入结果)。
5. **放不下** → 按 `daysStep` 切段**顺序循环**(单机 2 核 mr 并行收益小
且峰值内存翻倍):每段查询 `[段起点−L, 段终点+R]`(halo 重叠)、段内
截断,`mergeSegmentResults` 用 `concatMatrix(mats, false)` 纵向拼接并
重挂行/列标签。列已按 presentCodes 对齐,分段边界处滚动结果与整段
计算逐值一致。

### 分段占位统一为 NULL 矩阵(审查修复)

失败/空分段的占位由 `buildPlaceholderMatrix` 产出「行-该段交易日、列-
presentCodes」的 NULL DOUBLE 矩阵,与正常结果**同为带标签矩阵**。这样:

- `mergeSegmentResults` 对所有分段一视同仁按行拼接,**任一分段计算失败
不再丢弃其它分段的正确结果**(旧实现遇到占位 TABLE 就 `break`,会把整个
表达式的多年结果塌缩为单个短占位);
- 空 `baseData` 的分段也填占位而非整体缺失,**消除分段日期空洞**;
- 占位填 NULL(非 0),避免用 0 污染缺失位置的因子值;也消除了 Python 侧
`_computed_dict_to_panel` 期望三元素矩阵、却收到 TABLE 的类型隐患。

## 配置

| 键 | 默认 | 说明 |
|----|------|------|
| `ddb_lookback_default` | 252 | 回看解析兜底交易日数(仅对 qlib 解析不了且扫不到窗口整数的表达式生效) |
| `ddb_days_step` | 252 | 分段模式的段长(交易日);仅内存估算不足时启用分段 |

注意:`L`(回看天数)接近 `daysStep` 时分段 halo 开销接近 100%;表达式
含长窗口(如 252 日滚动)且触发分段时,可适当调大 `ddb_days_step`。

### 兜底窗口上界(审查修复)

回看正则兜底仅对 qlib 解析不了的表达式生效,取最大独立整数为窗口。为避免把
数值常量误判为窗口(如 `gtjaAlpha191_005($volume/1000000, ...)` 中的
`1000000` 被当成百万日回看,进而 `shiftTradeDays` 外扩到不合理区间、在
2 核/8GB 上触发超量查询),扫到的整数超过 `_MAX_FALLBACK_WINDOW`(2000,
约 8 年交易日)即视为常量剔除;剔除后无候选则退回 `ddb_lookback_default`。

## 验证状态

- 离线单测:`tests/test_ddb_lookback.py`(回看解析含大常量剔除/上界共 15 例
+ 脚本接线 3 例),`tests/test_fetch_features.py` 分支回归全部通过。
- **待 live 验证**(本机无 DDB 服务器):
1. `featureEngineering.dos` 语法与 `panel(...,, colLabels)` / `loc` /
`concatMatrix` / `rename!` / `matrix(DOUBLE,r,c,,double(NULL))` 行为;
2. 等价性:同一表达式(如 `Mean($close,20)`)文件后端 vs DDB 后端,
区间头部 L 天结果对齐;
3. 分段模式:临时把 `isRunWithAvailableMemory` 强制返回 false,比对
分段结果与整段结果逐值一致;并构造「某分段该表达式失败」用例,
确认其它分段结果不被丢弃、占位段为 NaN 且日期覆盖连续。
12 changes: 12 additions & 0 deletions qlib/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,18 @@ def register_from_C(config, skip_register=True):
"provider": "LocalProvider",
# database uri for database backend
"database_uri": None,
# ---- DolphinDB 后端批次参数(目标服务器为社区版 2 核/8GB 时的取舍)----
# 每批查询的字段数:增大会加宽客户端 concat 宽度、减少批次数
"ddb_field_chunk_size": 30,
# FeatureEngineeringByDate 的日期分片(交易日数,内存不足才启用分段):
# 单段服务器内存 ≈ days_step × 股票数 × 基础字段数 × 8B
# (5000 股 × 5 字段 × 252 天 ≈ 50MB),8GB 内存建议不超过 504
"ddb_days_step": 252,
# 表达式回看窗口的兜底交易日数:qlib 算子树解析不了(DDB 专属 alpha
# 库函数等)且正则扫不到窗口整数时,向前外扩这么多个交易日
"ddb_lookback_default": 252,
# True 时 init 全量预载三个 alpha 因子库(默认按字段引用惰性加载)
"ddb_preload_alpha_libs": False,
# dolphindb provider config
# "dolphindb_provider": None,
"dbclient_provider":None,
Expand Down
33 changes: 33 additions & 0 deletions qlib/data/backend/ddb_qlib/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -189,6 +189,39 @@ bridge.close()
- 做好异常处理和错误日志
- 定期检查数据一致性

5. **进程内缓存与失效**

为减少 RPC 往返,本后端在进程内缓存:交易日历(`TradeDateUtils`)、
表 schema(`utils.get_table_columns`)、`existsTable` 正向结果、
股票池(`H["i"]`)与表达式翻译(lru_cache)。

- 经 `write_df_to_ddb` / CSV 导入 / `clean_qlib_db` 的写入会**自动失效**缓存;
- 若数据被**外部进程**写入(本进程无感知),长驻只读进程需手动失效:

```python
from qlib.data.backend.ddb_qlib import invalidate_ddb_caches
invalidate_ddb_caches()
```

6. **批次参数(社区版 2 核/8GB 调优)**

```python
qlib.init(
database_uri="dolphindb://...",
ddb_field_chunk_size=30, # 每批查询的字段数
ddb_days_step=252, # 内存不足时的日期分段长(交易日数),8GB 建议 ≤504
ddb_lookback_default=252, # 表达式回看解析失败时的兜底外扩交易日数
ddb_preload_alpha_libs=False, # True 恢复 init 全量预载 alpha 因子库
)
```

计算分支自动按表达式外扩查询窗口(取前序期)并在服务器端截断回请求
区间;内存估算不足时按 `ddb_days_step` 分段顺序计算(每段带回看重叠),
详见 `docs/DDB取前序期与分段计算_20260722.md`。

可用 `DDB_BENCH_URI=... python scripts/benchmark_ddb_backend.py` 对
实际服务器做基准校准。

## 技术支持

如有问题,请联系:
Expand Down
34 changes: 33 additions & 1 deletion qlib/data/backend/ddb_qlib/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,41 @@
from .schemas import QlibTableSchema
from .ddb_features import (
register_ddb_functions_to_qlib,
# ddb_compute_features,
fetch_features_from_ddb,
normalize_fields_to_ddb,
TradeDateUtils,
adapt_qlib_expr_syntax_for_ddb,
)


def invalidate_ddb_caches() -> None:
"""清空 DDB 后端的进程内缓存。

覆盖范围:
- :class:`TradeDateUtils` 的模块级交易日历缓存;
- 表列名缓存(:func:`utils.get_table_columns`);
- 存储层 existsTable 正向缓存(``DBStorageMixin._exists_cache``);
- ``H["c"]`` 原始日历缓存与 ``H["i"]`` 股票池缓存。

写路径(``write_df_to_ddb`` / CSV 导入 / ``clean_qlib_db``)变更表数据后
会自动调用;长驻只读进程若感知到外部写入,也可手动调用本函数。
"""
from .utils import clear_table_columns_cache

TradeDateUtils.clear_cache()
clear_table_columns_cache()
try:
from qlib.data.storage.dolphindb_storage import DBStorageMixin

DBStorageMixin._exists_cache.clear()
except Exception: # pragma: no cover - storage 未加载时跳过
pass
try:
from qlib.data.cache import H

for unit_key in ("c", "i"):
unit = H[unit_key]
while len(unit):
unit.popitem(last=False)
except Exception: # pragma: no cover - 缓存清理失败不应阻断写路径
pass
Loading