主题切换
17.3 数据质量监控
名词解释:量化数据质量(Quantitative Data Quality)
数据质量是量化策略绩效的隐性决定因素。"垃圾进,垃圾出"(Garbage In, Garbage Out, GIGO)在量化领域带来的后果不仅是错误的研究结论,更可能是直接的交易损失。数据质量监控体系确保数据在完整性、准确性、一致性和及时性四个维度上满足研究和交易的要求。
一、数据完整性检查
1.1 完整性检查的维度
数据完整性检查覆盖以下维度:
| 检查维度 | 内容 | 典型问题 | 检查频率 |
|---|---|---|---|
| 缺失检查 | 应该存在的数据是否存在 | 某日缺少数据文件、某时段数据缺失 | 每日 |
| 重复检查 | 是否存在完全重复的记录 | ETL 重跑导致数据重复 | 实时/每日 |
| 覆盖度检查 | 标的覆盖度是否完整 | 某股票退市后数据中断 | 每月 |
| 字段完整性 | 关键字段是否为 NULL | 价格为空、成交量缺失 | 实时 |
1.2 完整性检查的实现
此处有展示代码展开 ▼
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 异常值类型
金融数据的异常值分为三类:
- 价格异常:相对于前一笔价格的跳变超过合理范围
- 成交量异常:单笔成交量远超历史分布
- 关系异常:违反价格关系(如 Bid > Ask, High < Low)
2.2 多级异常检测
此处有展示代码展开 ▼
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 自动化修复策略
此处有展示代码展开 ▼
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):经过多次修订后的版本
此处有展示代码展开 ▼
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 多源对账
此处有展示代码展开 ▼
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`(跨数据源对账:比较两个数据源的一致性)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
数据质量监控是量化基础设施中投资回报率最高的投入之一——早发现一个数据问题,可能避免一个基于错误数据的策略上线。当数据质量有了保障,下一个架构性问题是如何选择处理范式:批处理还是流处理。