📑 本文目录
数据仓库治理DuckDBschema合约字段缺失检测数据质量门禁可转债数据工程ETL校验

可转债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;
字段名类型类别
codeVARCHAR标识
dateVARCHAR标识
openDOUBLEOHLCV
highDOUBLEOHLCV
lowDOUBLEOHLCV
closeDOUBLEOHLCV
volDOUBLEOHLCV
amountDOUBLEOHLCV

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_pathfile_sizerow_counterror_flagupdated_at
/home/user/quant/data/bonds/110002.csv17758904321426002026-07-10T16:41
/home/user/quant/data/bonds/110003.csv17758897925196402026-07-10T16:41
/home/user/quant/data/bonds/110004.csv17758895361070802026-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_dailyOHLCV + 成交量/额✅ 已入库
条款层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_dailyOHLCV 齐全,覆盖基本分析需求
🟡 部分bond_daily, concept_daily, industry_daily有行情数据,缺衍生指标
🔴 受限stock_money_flow(曾消失过)、stock_valuation(曾停更38天)数据质量不稳定,需重点监控

将这些表全部纳入 schema_contracts.yaml,就能构建一个覆盖全库的数据质量监控体系。我在 数据缺口决策框架 中提过”采/买/算/代理”四象限——schema 合约就是那道”统一的管理界面”,不管字段是采集的、购买的还是计算的,都得先过门禁。

八、局限性说明

这套方案不是万能的,有几个实际的局限:

  1. YTM 补算依赖条款数据质量bond_basic 表的到期日、票面利率等数据本身可能有错误或滞后。垃圾进、垃圾出——衍生计算的正确性上限取决于源数据质量。

  2. schema 合约是静态的:每次新增字段需求需要手动更新 YAML。未来可以考虑从策略代码的 SELECT 语句中自动提取字段依赖,但工程复杂度较高。

  3. 门禁只检查”有/没有”,不检查”对不对”:字段存在但值全为 0 或异常值,schema 门禁不会发现。需要补充数据值域校验(如 close 大于 0、vol 非负等规则)。

  4. 性能开销:全库审计在 30 张表、75 万行+ 的规模下耗时可控(秒级),但随着数据增长,需要考虑分区校验或采样校验。

  5. 历史数据补算成本:4615 个交易日的 YTM 回算需要逐日执行现金流贴现,单线程全量补算可能需要数小时。实际操作建议只补算最近 2 年的数据,更早的按需计算。

九、总结

一条 SELECT ytm FROM bond_daily 的报错,暴露的不是单个字段的问题,而是数据仓库治理体系的缺位。

我给出的三层治理方案:

层次做什么工具
Schema 合约定义每张表的字段最低标准YAML + DuckDB DESCRIBE
衍生补算管道从已有数据计算缺失的衍生字段Python + AKShare + 牛顿法
数据质量门禁ETL 后自动校验,失败即告警cron + 飞书 webhook

这套体系的核心原则是:字段缺失应该是一个被主动检测的违规事件,而不是一个被被动发现的运行时错误。 当 cron 告诉你”bond_daily 缺 ytm”,你可以在分析策略时就知道该用什么代理指标;而当你在运行策略时才报错发现字段不存在,损失的是你的调试时间和策略可信度。

数据仓库治理不是一次性项目,而是一套持续运行的纪律。从 bond_daily 开始,逐步覆盖全库。


相关阅读

💬 评论