Skip to content

17.3 数据质量监控

名词解释:量化数据质量(Quantitative Data Quality)

数据质量是量化策略绩效的隐性决定因素。"垃圾进,垃圾出"(Garbage In, Garbage Out, GIGO)在量化领域带来的后果不仅是错误的研究结论,更可能是直接的交易损失。数据质量监控体系确保数据在完整性、准确性、一致性和及时性四个维度上满足研究和交易的要求。

一、数据完整性检查

1.1 完整性检查的维度

数据完整性检查覆盖以下维度:

检查维度内容典型问题检查频率
缺失检查应该存在的数据是否存在某日缺少数据文件、某时段数据缺失每日
重复检查是否存在完全重复的记录ETL 重跑导致数据重复实时/每日
覆盖度检查标的覆盖度是否完整某股票退市后数据中断每月
字段完整性关键字段是否为 NULL价格为空、成交量缺失实时

1.2 完整性检查的实现

PYTHON112 行 · 3.8 KB
📄此处有展示代码112 行 · 3.8 KB展开 ▼
python
import numpy as np
import pandas as pd
from datetime import datetime, timedelta
from typing import Dict, List, Tuple
import logging

class DataCompletenessChecker:
    """
    数据完整性检查器。

    针对 Tick 数据和日频数据的完整性检查。
    """

    def __init__(self, expected_symbols: list, expected_fields: list):
        self.expected_symbols = set(expected_symbols)
        self.expected_fields = expected_fields
        self.issues = []

    def check_missing_symbols(self,
                               data: pd.DataFrame,
                               date: str) -> Dict[str, List[str]]:
        """
        检查指定日期是否有 symbol 缺失。

        返回:
            {'missing': [...], 'extra': [...]}
        """
        present_symbols = set(data['symbol'].unique())
        missing = self.expected_symbols - present_symbols
        extra = present_symbols - self.expected_symbols

        if missing:
            self.issues.append({
                'type': 'missing_symbols',
                'date': date,
                'symbols': list(missing),
                'severity': 'HIGH'
            })

        return {'missing': list(missing), 'extra': list(extra)}

    def check_time_continuity(self,
                               data: pd.DataFrame,
                               symbol: str,
                               market_open: str,
                               market_close: str,
                               max_gap_minutes: int = 30) -> List[dict]:
        """
        检查某标的在交易时段内是否有超过阈值的间断。

        参数:
            data: 特定 symbol 的数据,需包含 timestamp 列
            symbol: 标的代码
            market_open, market_close: 交易时段(如 '09:30', '16:00')
            max_gap_minutes: 最大允许间断(分钟)
        """
        data = data.sort_values('timestamp')
        gaps = []
        max_gap = pd.Timedelta(minutes=max_gap_minutes)

        # 检查开盘是否有数据
        open_time = pd.Timestamp(f"{data['timestamp'].iloc[0].date()} {market_open}")
        first_record = data['timestamp'].iloc[0]

        if first_record - open_time > max_gap:
            gaps.append({
                'type': 'late_open',
                'symbol': symbol,
                'expected_open': str(open_time),
                'first_record': str(first_record),
                'gap_seconds': (first_record - open_time).total_seconds()
            })

        # 检查期间是否有间断
        time_diffs = data['timestamp'].diff()
        large_gaps = time_diffs[time_diffs > max_gap]

        for idx in large_gaps.index:
            gaps.append({
                'type': 'data_gap',
                'symbol': symbol,
                'gap_start': str(data.loc[idx, 'timestamp']),
                'gap_duration_seconds': large_gaps[idx].total_seconds(),
                'severity': 'HIGH' if large_gaps[idx].total_seconds() > 3600
                            else 'MEDIUM'
            })

        return gaps

    def check_field_completeness(self,
                                  data: pd.DataFrame) -> Dict[str, float]:
        """
        检查各字段的完整率(非空比例)。
        """
        completeness = {}
        total = len(data)

        for field in self.expected_fields:
            if field in data.columns:
                non_null = data[field].notna().sum()
                completeness[field] = non_null / total * 100

                if completeness[field] < 99.0:
                    self.issues.append({
                        'type': 'field_incompleteness',
                        'field': field,
                        'completeness_pct': completeness[field],
                        'severity': 'HIGH' if completeness[field] < 95.0
                                    else 'MEDIUM'
                    })

        return completeness
点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `DataCompletenessChecker`(数据完整性检查器)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。

二、异常值检测与自动化修复

2.1 异常值类型

金融数据的异常值分为三类:

  1. 价格异常:相对于前一笔价格的跳变超过合理范围
  2. 成交量异常:单笔成交量远超历史分布
  3. 关系异常:违反价格关系(如 Bid > Ask, High < Low)

2.2 多级异常检测

PYTHON73 行 · 2.5 KB
📄此处有展示代码73 行 · 2.5 KB展开 ▼
python
class AnomalyDetector:
    """
    多级异常值检测器。

    检测策略:
    Level 1: 静态规则(如价格 > 0, Ask > Bid)
    Level 2: 统计规则(如 Z-score, MAD-based)
    Level 3: 时序规则(如前后价格变动过大)
    """

    @staticmethod
    def level1_static_rules(row: pd.Series) -> List[str]:
        """一级:静态合理性规则"""
        violations = []

        if row.get('price', 0) <= 0:
            violations.append('PRICE_NEGATIVE_OR_ZERO')
        if row.get('size', 0) <= 0:
            violations.append('SIZE_NEGATIVE_OR_ZERO')
        if row.get('ask_price', 0) < row.get('bid_price', 0):
            violations.append('ASK_LESS_THAN_BID')
        if row.get('high', 0) < row.get('low', 0):
            violations.append('HIGH_LESS_THAN_LOW')
        if row.get('close', 0) > row.get('high', 0) or \
           row.get('close', 0) < row.get('low', 0):
            violations.append('CLOSE_OUTSIDE_HL_RANGE')

        return violations

    @staticmethod
    def level2_statistical_rules(data: pd.Series,
                                  method: str = 'iqr',
                                  multiplier: float = 5.0) -> np.ndarray:
        """
        二级:基于统计分布的异常检测。

        参数:
            data: 数值序列
            method: 'iqr' (四分位距) 或 'mad' (中位数绝对偏差)
            multiplier: 判定阈值乘数
        返回:
            布尔数组,True 表示异常
        """
        if method == 'iqr':
            q1 = data.quantile(0.25)
            q3 = data.quantile(0.75)
            iqr = q3 - q1
            lower = q1 - multiplier * iqr
            upper = q3 + multiplier * iqr
            return (data < lower) | (data > upper)

        elif method == 'mad':
            median = data.median()
            mad = np.median(np.abs(data - median))
            modified_z = 0.6745 * (data - median) / (mad + 1e-10)
            return np.abs(modified_z) > multiplier

    @staticmethod
    def level3_price_jump(data: pd.DataFrame,
                           max_bps_change: float = 500) -> np.ndarray:
        """
        三级:价格跳变检测。

        检测两笔连续交易间价格的异常跳变(bps)。

        参数:
            data: 含 price 列的时序数据
            max_bps_change: 最大允许的价格变动(基点)
        返回:
            异常标志
        """
        price_change_bps = data['price'].pct_change().abs() * 10000
        return price_change_bps > max_bps_change
点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `AnomalyDetector`(多级异常值检测器)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。

2.3 自动化修复策略

PYTHON78 行 · 2.9 KB
📄此处有展示代码78 行 · 2.9 KB展开 ▼
python
class AutoRepair:
    """
    自动修复常见数据质量问题。

    修复策略分层:
    1. 可直接修复的问题(如数据格式)
    2. 可推断修复的问题(如用前后值插补)
    3. 需标记后人工审核的问题
    """

    @staticmethod
    def repair_duplicate_timestamps(data: pd.DataFrame,
                                     timestamp_col: str = 'timestamp',
                                     strategy: str = 'keep_last') -> pd.DataFrame:
        """
        修复重复时间戳。

        策略:
        - 'keep_last': 保留最后一次出现的记录
        - 'keep_first': 保留第一次
        - 'average_price': 对成交价取成交量加权平均
        """
        if strategy == 'keep_last':
            return data.drop_duplicates(subset=[timestamp_col], keep='last')
        elif strategy == 'keep_first':
            return data.drop_duplicates(subset=[timestamp_col], keep='first')
        elif strategy == 'average_price':
            # 成交量加权平均
            grouped = data.groupby(timestamp_col)
            repaired = grouped.agg({
                'price': lambda x: np.average(x, weights=data.loc[x.index, 'size']),
                'size': 'sum'
            }).reset_index()
            return repaired

    @staticmethod
    def interpolate_small_gaps(data: pd.DataFrame,
                                time_col: str = 'timestamp',
                                value_col: str = 'price',
                                max_gap_seconds: int = 5) -> pd.DataFrame:
        """
        对短时间间断进行线性插值修复。

        只修复 max_gap_seconds 内的短间断,
        长间断保留为空(标记为质量问题)。
        """
        data = data.sort_values(time_col).reset_index(drop=True)
        data['time_diff'] = data[time_col].diff().dt.total_seconds()

        # 识别短间断
        short_gaps = (data['time_diff'] > 0) & \
                     (data['time_diff'] <= max_gap_seconds)

        if short_gaps.any():
            data[value_col] = data[value_col].interpolate(
                method='linear',
                limit=1  # 每个间断只补一个点
            )
            data['interpolated'] = short_gaps

        return data

    @staticmethod
    def flag_for_review(data: pd.DataFrame,
                         anomaly_mask: np.ndarray,
                         issue_type: str) -> pd.DataFrame:
        """
        标记异常记录,添加质量标记字段。

        通过添加 data_quality_flag 和 review_required 字段,
        保留数据供后续分析的同时标记其可信度。
        """
        data = data.copy()
        data['data_quality_flag'] = data.get('data_quality_flag', 'CLEAN')
        data.loc[anomaly_mask, 'data_quality_flag'] = issue_type
        data['review_required'] = data['data_quality_flag'] != 'CLEAN'

        return data
点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `AutoRepair`(自动修复常见数据质量问题)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。

三、数据版本管理与回溯

3.1 数据版本管理的必要性

金融数据具有"后见之明"特征:发布时间表、修订和修正使得数据的"最新版本"与历史时点的"当时可得版本"可能不同:

  • 实时版本(Real-time):数据发布时的初始版本
  • 修订版本(Revised):后续修正的版本
  • 最终版本(Final):经过多次修订后的版本
PYTHON55 行 · 1.9 KB
📄此处有展示代码55 行 · 1.9 KB展开 ▼
python
class DataVersionControl:
    """
    量化数据的版本控制。

    使用时间旅行查询(Time Travel)获取任意历史时点的
    "当时已知"的数据快照,防止前视偏差(Look-ahead Bias)。
    """

    def __init__(self, storage_backend):
        self.backend = storage_backend

    def write_versioned(self,
                         data: pd.DataFrame,
                         dataset: str,
                         as_of_date: str,
                         is_revision: bool = False):
        """
        写入带版本的数据。

        参数:
            data: 数据
            dataset: 数据集名称(如 'financials', 'estimates')
            as_of_date: 数据发布日
            is_revision: 是否为对已发布数据的修订
        """
        version = {
            'dataset': dataset,
            'as_of_date': as_of_date,
            'ingestion_timestamp': datetime.now().isoformat(),
            'is_revision': is_revision,
            'version_number': self._get_next_version(dataset, as_of_date)
        }

        self.backend.write(data, metadata=version)

    def get_as_of(self,
                   dataset: str,
                   observation_date: str,
                   as_of_date: str) -> pd.DataFrame:
        """
        获取 'as_of_date' 时点已知的、关于 'observation_date' 的数据。

        这禁止"使用未来信息"的 look-ahead bias。
        """
        # 查询 observation_date 之前、as_of_date 之前的最新版本
        return self.backend.query(
            dataset=dataset,
            observation_date=observation_date,
            ingested_before=as_of_date
        )

    def _get_next_version(self, dataset: str, as_of_date: str) -> int:
        """获取下一个版本号"""
        existing = self.backend.list_versions(dataset, as_of_date)
        return len(existing) + 1
点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `DataVersionControl`(量化数据的版本控制)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。

四、对账机制

4.1 多源对账

PYTHON69 行 · 2.3 KB
📄此处有展示代码69 行 · 2.3 KB展开 ▼
python
def cross_source_reconciliation(source_a: pd.DataFrame,
                                 source_b: pd.DataFrame,
                                 key_columns: list,
                                 value_columns: list,
                                 tolerance: dict = None) -> dict:
    """
    跨数据源对账:比较两个数据源的一致性。

    参数:
        source_a, source_b: 两个数据源的 DataFrames
        key_columns: 连接键(如 ['symbol', 'date'])
        value_columns: 要比较的数值列
        tolerance: {列名: 容忍度},如 {'price': 0.01, 'volume': 100}
    返回:
        对账结果,包含差异统计和详细差异记录
    """
    if tolerance is None:
        tolerance = {}

    # 内连接
    merged = source_a.merge(
        source_b,
        on=key_columns,
        how='outer',
        suffixes=('_A', '_B'),
        indicator=True
    )

    # 只在A/B中存在的记录
    only_in_a = merged[merged['_merge'] == 'left_only']
    only_in_b = merged[merged['_merge'] == 'right_only']
    in_both = merged[merged['_merge'] == 'both']

    # 比较值列
    differences = []
    for col in value_columns:
        col_a = f"{col}_A"
        col_b = f"{col}_B"

        if col_a in in_both.columns and col_b in in_both.columns:
            diff = (in_both[col_a] - in_both[col_b]).abs()
            tol = tolerance.get(col, 0)

            mismatch = diff > tol
            n_mismatch = mismatch.sum()

            if n_mismatch > 0:
                max_diff = diff.max()
                mean_diff = diff[mismatch].mean()
                differences.append({
                    'column': col,
                    'n_mismatches': n_mismatch,
                    'mismatch_rate': n_mismatch / len(in_both) * 100,
                    'max_absolute_diff': max_diff,
                    'mean_absolute_diff': mean_diff
                })

    return {
        'total_records_A': len(source_a),
        'total_records_B': len(source_b),
        'matched_records': len(in_both),
        'only_in_A': len(only_in_a),
        'only_in_B': len(only_in_b),
        'match_rate': len(in_both) / max(len(source_a), len(source_b)) * 100,
        'differences': differences,
        'passes_reconciliation': len(differences) == 0 and
                                  len(only_in_a) == 0 and
                                  len(only_in_b) == 0
    }
点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:函数 `cross_source_reconciliation`(跨数据源对账:比较两个数据源的一致性)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。

数据质量监控是量化基础设施中投资回报率最高的投入之一——早发现一个数据问题,可能避免一个基于错误数据的策略上线。当数据质量有了保障,下一个架构性问题是如何选择处理范式:批处理还是流处理