📑 本文目录
可转债YTM排名跑不出来?一张缺失字段暴露的量化数据仓库 schema 治理清单
引子:一条 SELECT 语句引发的治理思考
2026年7月13日,我在执行每日信息采集 cron 时,尝试跑一个最基础的可转债分析——到期收益率(YTM)排名,找出估值最低的可转债 top 10。SQL 很简单:
SELECT code, close, ytm
FROM bond_daily
WHERE date = '2026-07-10'
ORDER BY ytm ASC
LIMIT 10;
结果 DuckDB 直接报错:
Catalog Error: Column "ytm" does not exist in table "bond_daily"
这张表有 75 万行数据、覆盖 988 只可转债、横跨 2007 年至今 4615 个交易日——数据量庞大,但 ytm 字段根本不存在。一张有数据的表,却跑不出一个最基础的量化分析。
这不是第一次遇到这个问题。我之前写过一篇 可转债 YTM 计算管道 从算法层解决了”如何补算 YTM”。但这次我停下来想了一个更深层的问题:
为什么这个问题反复出现? 不是因为不会算 YTM,而是因为整个数据仓库没有字段级完整性治理。ETL 脚本只拉了 OHLCV 就算完事,没有任何机制在写入前检查”这张表是否满足量化分析所需的最低字段要求”。
本文就来补齐这块治理短板。
一、问题拆解:bond_daily 表结构审计
先看清问题全貌。用 DuckDB 的 DESCRIBE 审计 bond_daily 表:
DESCRIBE bond_daily;
| 字段名 | 类型 | 类别 |
|---|---|---|
code | VARCHAR | 标识 |
date | VARCHAR | 标识 |
open | DOUBLE | OHLCV |
high | DOUBLE | OHLCV |
low | DOUBLE | OHLCV |
close | DOUBLE | OHLCV |
vol | DOUBLE | OHLCV |
amount | DOUBLE | OHLCV |
8 个字段,全部是 OHLCV 行情数据。 对比一下可转债量化策略实际需要的字段:
| 策略类型 | 所需字段 | bond_daily 有没有 |
|---|---|---|
| 双低策略 | close, 转股溢价率 | close ✅ 溢价率 ❌ |
| 到期收益率排序 | close, ytm, 剩余年限 | close ✅ ytm ❌ 年限 ❌ |
| 下修博弈 | 是否触发改下修条件, 转股价 | 全部 ❌ |
| 隐含波动率 | close, 正股价格, 转股价, 波动率 | close ✅ 其余 ❌ |
| 条款触发监控 | 回售价, 赎回价, 转股价 | 全部 ❌ |
9 个关键字段中有 8 个缺失。 这不是某几个策略不能用的问题,而是整个可转债分析体系的基础设施缺了一大半。
更值得关注的是——数据库里甚至没有 bond_basic 表来存储可转债的基础信息(到期日、转股价、票面利率等)。我查了全部 30 张表:
SELECT table_name FROM information_schema.tables
WHERE table_schema = 'main' AND table_name LIKE '%bond%';
结果只有 bond_daily(行情)和 bond_yield_curve(国债收益率曲线),可转债条款信息表完全不存在。
二、影响面评估:缺字段连锁波及了什么
字段缺失不是孤立问题,它会沿依赖链向上传导。我追踪了 AgentQuant 博客中所有涉及可转债的文章:
| 文章 | 依赖字段 | 受影响程度 |
|---|---|---|
| 可转债18策略横评 | 双低、YTM、溢价率 | 高 |
| 可转债YTM计算管道 | 专门解决YTM补算 | 本文前置 |
| 双低策略失效分析 | close, 溢价率 | 中 |
| 低价轮动策略 | close | 低(仅用价格) |
| 网格交易策略 | high, low, close | 低(仅用OHLCV) |
连锁反应链路:bond_daily 缺字段 → 双低排名跑不出 → 18 策略横评只能用代理指标 → 策略回测结论可信度下降。
我在 数据质量五大噩梦 中提过,可转债数据停更 24 天是时效性问题。但字段缺失是比数据停更更隐蔽的问题——数据停更你会立刻发现(查不到当天数据),而字段缺失是你以为有数据、实际上跑不出有用的分析,这种”静默失效”更危险。
三、根因分析:ETL 字段映射断层
为什么会缺这么多字段?我查了 ETL 检查点表 _etl_checkpoint,追踪数据入库逻辑:
SELECT * FROM _etl_checkpoint
WHERE file_path LIKE '%bonds%'
ORDER BY updated_at DESC
LIMIT 5;
| file_path | file_size | row_count | error_flag | updated_at |
|---|---|---|---|---|
| /home/user/quant/data/bonds/110002.csv | 1775890432 | 14260 | 0 | 2026-07-10T16:41 |
| /home/user/quant/data/bonds/110003.csv | 1775889792 | 51964 | 0 | 2026-07-10T16:41 |
| /home/user/quant/data/bonds/110004.csv | 1775889536 | 10708 | 0 | 2026-07-10T16:41 |
关键发现:ETL 脚本从 CSV 文件读取数据后,只映射了 OHLCV 这 8 个字段就写入 DuckDB。 源数据 CSV 可能包含更多信息(取决于 AKShare 的 bond_zh_hs_daily 接口返回字段),但 ETL 的 INSERT 语句写死了 8 列映射:
# 当前 ETL 脚本的简化逻辑(问题根源)
con.execute("""
INSERT INTO bond_daily (code, date, open, high, low, close, vol, amount)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
""", row_data)
而 YTM、转股溢价率等字段在 AKShare 的行情接口中本来就不返回——它们是衍生计算字段,需要从条款数据(bond_zh_cov 接口)获取到期日、票面利率等基础信息后再计算。
这就是字段映射断层的典型模式:
| 层次 | 数据源 | 有哪些字段 | 断点 |
|---|---|---|---|
| 行情层 | AKShare bond_zh_hs_daily | OHLCV + 成交量/额 | ✅ 已入库 |
| 条款层 | AKShare bond_zh_cov | 到期日、转股价、票面利率、回售条件 | ❌ 未入库 |
| 衍生层 | 自行计算 | YTM、溢价率、剩余年限、delta | ❌ 未计算 |
从行情到衍生,中间断了两层。ETL 脚本做了第一步,后面的全断了。
四、治理方案一:Schema 合约——用 SQL 定义”字段最低标准”
核心思路:每张表都必须声明它承诺提供的字段清单。 ETL 写入前,自动校验源数据是否满足这份合约;不满足就拒绝写入并发告警。
我用一个 YAML 文件来定义 schema 合约,它清晰、可版本控制、非技术人员也能读懂:
# schema_contracts.yaml
# 定义每张表必须满足的字段清单
bond_daily:
description: "可转债日线行情"
required_fields:
- name: code
type: VARCHAR
category: 标识
- name: date
type: VARCHAR
category: 标识
- name: open
type: DOUBLE
category: OHLCV
- name: high
type: DOUBLE
category: OHLCV
- name: low
type: DOUBLE
category: OHLCV
- name: close
type: DOUBLE
category: OHLCV
- name: vol
type: DOUBLE
category: OHLCV
- name: amount
type: DOUBLE
category: OHLCV
# 期望但非阻塞的字段(缺失时告警但不阻止写入)
expected_fields:
- name: ytm
type: DOUBLE
category: 衍生指标
reason: "双低策略/到期收益率排序必需"
- name: premium_rate
type: DOUBLE
category: 衍生指标
reason: "转股溢价率,几乎所有可转债策略都依赖"
- name: remain_year
type: DOUBLE
category: 衍生指标
reason: "剩余年限,用于到期收益率计算"
quality_rules:
- rule: "close 大于 0"
sql: "COUNT(*) WHERE close <= 0 = 0"
- rule: "每日记录数在 300-1200 之间"
sql: "daily_count BETWEEN 300 AND 1200"
注意我区分了 required_fields(必须满足,缺失即拒绝写入)和 expected_fields(期望有,缺失时告警但允许继续)。这个设计很重要——YTM 等衍生字段需要额外计算管道才能产出,强行要求会导致整个 ETL 卡住;但告警能确保你知道它在缺。
对应的校验脚本:
import duckdb
import yaml
from pathlib import Path
def validate_schema(con, table_name, contract):
"""校验表是否满足 schema 合约,返回违规清单"""
violations = []
# 获取实际表结构
actual_cols = con.execute(f"DESCRIBE {table_name}").fetchall()
actual_col_names = {row[0]: row[1] for row in actual_cols}
# 检查 required_fields
for field in contract.get('required_fields', []):
if field['name'] not in actual_col_names:
violations.append({
'level': 'CRITICAL',
'table': table_name,
'field': field['name'],
'msg': f"必需字段缺失: {field['name']} ({field['type']})"
})
elif actual_col_names[field['name']] != field['type']:
violations.append({
'level': 'WARNING',
'table': table_name,
'field': field['name'],
'msg': f"类型不匹配: 期望 {field['type']}, 实际 {actual_col_names[field['name']]}"
})
# 检查 expected_fields
for field in contract.get('expected_fields', []):
if field['name'] not in actual_col_names:
violations.append({
'level': 'ALERT',
'table': table_name,
'field': field['name'],
'msg': f"期望字段缺失: {field['name']} — {field.get('reason', '')}"
})
return violations
def run_full_audit(db_path, contracts_path):
"""执行全库 schema 审计"""
con = duckdb.connect(db_path, read_only=True)
contracts = yaml.safe_load(Path(contracts_path).read_text())
all_violations = []
for table_name, contract in contracts.items():
try:
violations = validate_schema(con, table_name, contract)
all_violations.extend(violations)
except Exception as e:
all_violations.append({
'level': 'CRITICAL',
'table': table_name,
'field': 'N/A',
'msg': f"表不存在或无法访问: {e}"
})
return all_violations
# 执行审计
violations = run_full_audit(
'/mnt/c/Users/Administrator/clawd/data/quant/quant_v2.duckdb',
'schema_contracts.yaml'
)
for v in violations:
print(f"[{v['level']}] {v['table']}.{v['field']}: {v['msg']}")
跑一下当前 bond_daily 的审计结果:
[ALERT] bond_daily.ytm: 期望字段缺失: ytm — 双低策略/到期收益率排序必需
[ALERT] bond_daily.premium_rate: 期望字段缺失: premium_rate — 转股溢价率,几乎所有可转债策略都依赖
[ALERT] bond_daily.remain_year: 期望字段缺失: remain_year — 剩余年限,用于到期收益率计算
三条 ALERT,没有 CRITICAL——说明 OHLCV 基础字段齐全,但衍生指标层完全空白。这份报告就是治理的起点。
五、治理方案二:衍生字段计算管道——从已有数据补算缺失字段
知道了缺什么,下一步是补算。在 YTM 计算管道 一文中我已详细讲过 YTM 的牛顿法求解和转股溢价率计算,这里不重复算法细节,重点讲如何把它工程化为自动化补算管道。
核心思路是三步走:
第一步:建一张 bond_basic 表存储可转债条款信息
CREATE TABLE IF NOT EXISTS bond_basic (
code VARCHAR PRIMARY KEY,
name VARCHAR,
list_date VARCHAR,
delist_date VARCHAR,
maturity_date VARCHAR, -- 到期日
par_value DOUBLE, -- 面值(通常100)
coupon_rate VARCHAR, -- 票面利率(逐年递增,存JSON字符串)
convert_price DOUBLE, -- 转股价
sellback_price DOUBLE, -- 回售价
redeem_price DOUBLE, -- 赎回价
updated_at TIMESTAMP
);
这张表是衍生计算的”原料库”。数据源是 AKShare 的 bond_zh_cov 接口(可转债详情),每日全量刷新。
第二步:建一张 bond_metrics 表存储衍生指标
CREATE TABLE IF NOT EXISTS bond_metrics (
code VARCHAR,
date VARCHAR,
close DOUBLE,
ytm DOUBLE, -- 到期收益率
premium_rate DOUBLE, -- 转股溢价率
remain_year DOUBLE, -- 剩余年限
bond_premium DOUBLE, -- 纯债溢价率
delta DOUBLE, -- delta值(可选)
updated_at TIMESTAMP,
PRIMARY KEY (code, date)
);
第三步:每日补算管道
def daily_bond_metrics_pipeline(con, target_date):
"""每日衍生指标补算管道"""
# 1. 读取当日行情 + 条款信息
df = con.execute("""
SELECT d.code, d.date, d.close,
b.maturity_date, b.par_value, b.coupon_rate,
b.convert_price, b.sellback_price
FROM bond_daily d
JOIN bond_basic b ON d.code = b.code
WHERE d.date = ?
""", [target_date]).fetchdf()
if df.empty:
print(f"[WARN] {target_date} 无可转债行情数据")
return 0
# 2. 逐券计算衍生指标
results = []
for _, row in df.iterrows():
metrics = {
'code': row['code'],
'date': row['date'],
'close': row['close'],
}
# YTM 计算(调用已有管道)
if row['maturity_date'] and row['coupon_rate']:
metrics['ytm'] = calc_ytm(
current_price=row['close'],
maturity_date=row['maturity_date'],
coupon_schedule=parse_coupon(row['coupon_rate']),
par_value=row['par_value'] or 100,
eval_date=target_date
)
metrics['remain_year'] = years_between(target_date, row['maturity_date'])
# 转股溢价率
if row['convert_price']:
stock_price = get_underlying_price(con, row['code'], target_date)
if stock_price:
convert_value = stock_price / row['convert_price'] * 100
metrics['premium_rate'] = (row['close'] - convert_value) / convert_value * 100
results.append(metrics)
# 3. 写入 bond_metrics 表
if results:
con.register('temp_metrics', pd.DataFrame(results))
con.execute("""
INSERT OR REPLACE INTO bond_metrics
SELECT * FROM temp_metrics
""")
return len(results)
这里有几个工程细节值得注意:
INSERT OR REPLACE:DuckDB 支持 upsert 语义,补算重跑时不会产生重复数据- JOIN bond_basic:条款信息作为计算原料,缺条款的转债自然跳过(不报错,记录日志)
- 容错设计:
if row['convert_price']这类条件判断确保个别字段缺失不会让整个管道崩溃
六、落地:数据质量门禁 cron——让治理自动化
Schema 合约和补算管道写好了,但如果不自动化,很快又会回到”没人检查”的状态。我来构建一个每日数据质量门禁 cron,在 ETL 完成后自动执行。
门禁架构
每日 ETL 写入 bond_daily
↓
[门禁检查点] ← schema 合约校验
↓
通过 → 触发衍生指标补算
↓
[门禁检查点] ← 衍生指标完整性校验
↓
通过 → 生成当日可转债分析报告
↓
失败 → 飞书告警(附违规清单)
完整门禁脚本
#!/usr/bin/env python3.10
"""数据质量门禁 — 每日 ETL 后自动执行"""
import duckdb
import yaml
import json
from pathlib import Path
from datetime import datetime, date
DB_PATH = '/mnt/c/Users/Administrator/clawd/data/quant/quant_v2.duckdb'
CONTRACTS_PATH = '/home/user/quant/config/schema_contracts.yaml'
ALERT_LOG = '/tmp/dq_gate_alerts.json'
DAILY_REPORT_DIR = '/tmp/dq_reports'
def schema_gate(con, contracts):
"""第一道门禁:Schema 合约校验"""
violations = []
for table_name, contract in contracts.items():
try:
actual_cols = con.execute(f"DESCRIBE {table_name}").fetchall()
actual_names = {r[0] for r in actual_cols}
for field in contract.get('required_fields', []):
if field['name'] not in actual_names:
violations.append({
'gate': 'SCHEMA_REQUIRED',
'severity': 'CRITICAL',
'table': table_name,
'field': field['name'],
'msg': f"必需字段 {field['name']} 缺失"
})
for field in contract.get('expected_fields', []):
if field['name'] not in actual_names:
violations.append({
'gate': 'SCHEMA_EXPECTED',
'severity': 'ALERT',
'table': table_name,
'field': field['name'],
'msg': f"期望字段 {field['name']} 缺失: {field.get('reason','')}"
})
except Exception as e:
violations.append({
'gate': 'TABLE_ACCESS',
'severity': 'CRITICAL',
'table': table_name,
'msg': str(e)
})
return violations
def freshness_gate(con, table_name, date_col='date', max_lag_days=2):
"""第二道门禁:数据时效性检查"""
latest = con.execute(f"""
SELECT MAX({date_col}) FROM {table_name}
""").fetchone()[0]
today = date.today().isoformat()
lag = (datetime.fromisoformat(today) - datetime.fromisoformat(latest)).days
if lag > max_lag_days:
return [{
'gate': 'FRESHNESS',
'severity': 'WARNING',
'table': table_name,
'msg': f"最新数据日期 {latest},滞后 {lag} 天(阈值 {max_lag_days} 天)"
}]
return []
def completeness_gate(con, table_name, target_date):
"""第三道门禁:当日数据完整性检查"""
row_count = con.execute(f"""
SELECT COUNT(*) FROM {table_name} WHERE date = '{target_date}'
""").fetchone()[0]
# 可转债通常 300-700 只在交易
if row_count < 100:
return [{
'gate': 'COMPLETENESS',
'severity': 'CRITICAL',
'table': table_name,
'msg': f"{target_date} 仅 {row_count} 条记录,可能数据不全"
}]
return []
def run_gates():
"""执行全部门禁检查"""
con = duckdb.connect(DB_PATH, read_only=True)
contracts = yaml.safe_load(Path(CONTRACTS_PATH).read_text())
today = date.today().strftime('%Y-%m-%d')
all_violations = []
# Gate 1: Schema 合约
all_violations += schema_gate(con, contracts)
# Gate 2: 数据时效性
all_violations += freshness_gate(con, 'bond_daily', max_lag_days=2)
# Gate 3: 数据完整性
all_violations += completeness_gate(con, 'bond_daily', today)
# 分级处理
critical = [v for v in all_violations if v['severity'] == 'CRITICAL']
warnings = [v for v in all_violations if v['severity'] == 'WARNING']
alerts = [v for v in all_violations if v['severity'] == 'ALERT']
report = {
'date': today,
'total_violations': len(all_violations),
'critical': len(critical),
'warning': len(warnings),
'alert': len(alerts),
'details': all_violations
}
Path(ALERT_LOG).write_text(json.dumps(report, ensure_ascii=False, indent=2))
if critical:
print(f"🔴 {len(critical)} 个 CRITICAL 违规,ETL 可能失败,需人工介入")
for v in critical:
print(f" [{v['table']}] {v['msg']}")
return False # 门禁未通过
elif warnings or alerts:
print(f"🟡 {len(warnings)} 个 WARNING, {len(alerts)} 个 ALERT")
for v in warnings + alerts:
print(f" [{v['table']}] {v['msg']}")
return True # 门禁通过,但有告警
else:
print("✅ 所有数据质量门禁通过")
return True
if __name__ == '__main__':
passed = run_gates()
exit(0 if passed else 1)
挂 cron 每日自动执行
在 Hermes cron 中注册一个每日门禁任务:
# 每天 ETL 完成后(约 17:00)执行数据质量门禁
0 17 * * * /usr/bin/python3.10 /home/user/quant/scripts/dq_gate.py >> /tmp/dq_gate.log 2>&1
门禁失败时(返回非零退出码),配合飞书 webhook 推送告警:
# 在 run_gates() 的 critical 分支中添加
import requests
def send_feishu_alert(violations):
"""飞书机器人推送数据质量告警"""
msg = f"🔴 数据质量门禁告警 ({date.today()})\n"
for v in violations:
if v['severity'] == 'CRITICAL':
msg += f" • [{v['table']}] {v['msg']}\n"
requests.post(
'https://open.feishu.cn/open-apis/bot/v2/hook/<your-webhook>',
json={'msg_type': 'text', 'content': {'text': msg}}
)
七、扩展:从 bond_daily 到全库治理
可转债只是切入点。同样的 schema 合约 + 门禁模式可以推广到全库所有表。我对 quant_v2.duckdb 的 30 张表做了初步审计,按字段完整度分三档:
| 完整度 | 表 | 状态 |
|---|---|---|
| 🟢 完整 | stock_daily, index_daily, etf_daily, futures_daily | OHLCV 齐全,覆盖基本分析需求 |
| 🟡 部分 | bond_daily, concept_daily, industry_daily | 有行情数据,缺衍生指标 |
| 🔴 受限 | stock_money_flow(曾消失过)、stock_valuation(曾停更38天) | 数据质量不稳定,需重点监控 |
将这些表全部纳入 schema_contracts.yaml,就能构建一个覆盖全库的数据质量监控体系。我在 数据缺口决策框架 中提过”采/买/算/代理”四象限——schema 合约就是那道”统一的管理界面”,不管字段是采集的、购买的还是计算的,都得先过门禁。
八、局限性说明
这套方案不是万能的,有几个实际的局限:
-
YTM 补算依赖条款数据质量:
bond_basic表的到期日、票面利率等数据本身可能有错误或滞后。垃圾进、垃圾出——衍生计算的正确性上限取决于源数据质量。 -
schema 合约是静态的:每次新增字段需求需要手动更新 YAML。未来可以考虑从策略代码的
SELECT语句中自动提取字段依赖,但工程复杂度较高。 -
门禁只检查”有/没有”,不检查”对不对”:字段存在但值全为 0 或异常值,schema 门禁不会发现。需要补充数据值域校验(如 close 大于 0、vol 非负等规则)。
-
性能开销:全库审计在 30 张表、75 万行+ 的规模下耗时可控(秒级),但随着数据增长,需要考虑分区校验或采样校验。
-
历史数据补算成本:4615 个交易日的 YTM 回算需要逐日执行现金流贴现,单线程全量补算可能需要数小时。实际操作建议只补算最近 2 年的数据,更早的按需计算。
九、总结
一条 SELECT ytm FROM bond_daily 的报错,暴露的不是单个字段的问题,而是数据仓库治理体系的缺位。
我给出的三层治理方案:
| 层次 | 做什么 | 工具 |
|---|---|---|
| Schema 合约 | 定义每张表的字段最低标准 | YAML + DuckDB DESCRIBE |
| 衍生补算管道 | 从已有数据计算缺失的衍生字段 | Python + AKShare + 牛顿法 |
| 数据质量门禁 | ETL 后自动校验,失败即告警 | cron + 飞书 webhook |
这套体系的核心原则是:字段缺失应该是一个被主动检测的违规事件,而不是一个被被动发现的运行时错误。 当 cron 告诉你”bond_daily 缺 ytm”,你可以在分析策略时就知道该用什么代理指标;而当你在运行策略时才报错发现字段不存在,损失的是你的调试时间和策略可信度。
数据仓库治理不是一次性项目,而是一套持续运行的纪律。从 bond_daily 开始,逐步覆盖全库。
相关阅读
- 可转债YTM计算管道:从条款数据采集到每日刷新 — 算法层:如何计算 YTM 和转股溢价率
- 数据质量五大噩梦 — 时效性与故障模式:数据停更、ETL 滞后、cron 静默死亡
- 数据缺口决策框架 — 宏观层面:采/买/算/代理四象限
- 可转债18策略横评 — 策略层:哪些策略最受字段缺失影响
- DuckDB A股量化数据库搭建 — 数据库架构背景