Skip to content

Latest commit

 

History

History
1095 lines (898 loc) · 30.6 KB

File metadata and controls

1095 lines (898 loc) · 30.6 KB

WF-005: 历史数据清洗

📌 工作流基本信息

属性 内容
工作流ID WF-005
工作流名称 历史数据清洗
功能描述 对MatrixOne库中的历史Issue数据进行清洗、规范化、AI重新打标签,存入实验库用于优化和分析
实现状态 ❌ 未实现
云端可用 ✅ 设计上可用(纯数据库+AI操作)
核心价值 提升数据质量,为后续分析提供更准确的数据基础

🔄 流程步骤总览

┌──────────┐   ┌──────────┐   ┌──────────┐   ┌──────────┐   ┌──────────┐
│  步骤1   │──▶│  步骤2   │──▶│  步骤3   │──▶│  步骤4   │──▶│  步骤5   │
│读取历史  │   │数据清洗  │   │AI重新    │   │质量验证  │   │存入实验库│
│  Issue   │   │规范化    │   │ 打标签   │   │          │   │         │
└──────────┘   └──────────┘   └──────────┘   └──────────┘   └──────────┘
  从MO读取      去重/格式化     类型/优先级    检查完整性    experimental
  指定时间段     字段标准化     Labels规范     数据质量      _issues表
  + 过滤条件     + 数据补全     + AI校正       + 日志记录    + 清洗报告

快速理解

  1. 步骤1 - 从MO的issues_snapshot表读取历史Issue(可指定时间范围)
  2. 步骤2 - 数据清洗:去重、字段格式规范化、缺失值补全
  3. 步骤3 - AI重新分析:重新判断类型、优先级,规范化Labels
  4. 步骤4 - 质量验证:检查数据完整性、标签一致性
  5. 步骤5 - 存入实验库:写入experimental_issues表,生成清洗报告

核心特点:❌ 待实现 | ✅ 设计完整 | ✅ AI驱动清洗 | ✅ 可回溯


📥 整体输入

1. 数据源配置

输入项 类型 说明 示例
repo_owner String 仓库所有者 matrixorigin
repo_name String 仓库名称 matrixone
source_table String 源数据表 issues_snapshot

2. 时间范围配置

输入项 类型 说明 默认值
start_date Date 开始日期 None(全部)
end_date Date 结束日期 None(全部)
use_latest_snapshot Boolean 是否只用最新快照 True

3. 清洗规则配置

输入项 类型 说明 来源
cleaning_rules YAML 清洗规则配置文件 config/cleaning_rules.yaml
label_mapping Dict Labels标准化映射表 配置文件
field_validators Dict 字段验证规则 配置文件

cleaning_rules.yaml示例

# 去重规则
deduplication:
  method: "issue_id"  # 按issue_id去重
  keep: "latest"      # 保留最新的记录

# 字段规范化
field_normalization:
  title:
    max_length: 256
    trim: true
    remove_emoji: false
  
  state:
    valid_values: ["open", "closed"]
    default: "open"
  
  labels:
    format: "json_array"
    lowercase: false

# Labels标准化映射
label_mapping:
  # 统一命名
  "bug": "kind/bug"
  "feature": "kind/feature"
  "问数": "area/问数"
  "chatbi": "area/ChatBI"
  
  # 废弃标签移除
  deprecated:
    - "wontfix"
    - "invalid"

# 缺失值处理
missing_values:
  assignee:
    action: "keep_null"  # 保持NULL
  
  priority:
    action: "ai_infer"   # AI推断
  
  milestone:
    action: "set_default"
    default: "backlog"

4. AI服务配置

输入项 类型 说明
AI_PROVIDER String qwen(推荐)
DASHSCOPE_API_KEY String 通义千问API密钥
QWEN_MODEL String qwen-plus
batch_size Integer AI批量处理大小(默认50)

5. 质量检查配置

输入项 类型 说明 默认值
required_fields Array 必需字段列表 ["title", "state", "labels"]
label_validation Boolean 是否验证Labels有效性 True
min_quality_score Float 最低质量分数(0-1) 0.7

📤 整体输出

1. 清洗后的数据

输出项 表名 说明
清洗后Issue数据 experimental_issues 存储清洗和重新标注的Issue

2. 清洗报告

输出项 类型 说明
cleaning_report.json JSON 详细清洗报告
cleaning_report.md Markdown 可读性报告
data_quality_metrics Dict 数据质量指标

清洗报告内容

{
    "execution_info": {
        "start_time": "2026-03-04T10:00:00Z",
        "end_time": "2026-03-04T10:30:00Z",
        "duration_seconds": 1800,
        "repo": "matrixorigin/matrixone"
    },
    "data_statistics": {
        "total_issues_read": 5000,
        "duplicates_removed": 50,
        "records_cleaned": 4950,
        "records_failed": 5,
        "success_rate": 0.999
    },
    "cleaning_actions": {
        "field_normalization": 3200,
        "label_standardization": 2500,
        "ai_re_labeling": 4950,
        "missing_value_filled": 800
    },
    "quality_metrics": {
        "average_quality_score": 0.92,
        "issues_above_threshold": 4900,
        "issues_below_threshold": 50
    },
    "label_changes": {
        "total_labels_modified": 1500,
        "mapping_applied": 800,
        "ai_corrections": 700
    }
}

3. 变更日志

输出项 表名 说明
数据变更记录 data_cleaning_log 记录每条数据的修改详情

🔄 详细步骤拆分

步骤1: 读取历史Issue数据

步骤ID: WF-005-S01
功能: 从issues_snapshot表读取需要清洗的历史数据
实现状态: ❌ 待实现

输入

  • repo_owner, repo_name: 仓库标识
  • start_date, end_date: 时间范围(可选)
  • use_latest_snapshot: 是否只用最新快照

处理逻辑

读取策略1:最新快照(推荐)

-- 获取最新快照时间
SELECT MAX(snapshot_time) AS latest_time
FROM issues_snapshot
WHERE repo_owner = :owner AND repo_name = :repo;

-- 读取最新快照的所有Issue
SELECT *
FROM issues_snapshot
WHERE repo_owner = :owner 
  AND repo_name = :repo
  AND snapshot_time = :latest_time;

读取策略2:时间范围

SELECT *
FROM issues_snapshot
WHERE repo_owner = :owner 
  AND repo_name = :repo
  AND created_at >= :start_date
  AND created_at <= :end_date
  AND snapshot_time = (
      SELECT MAX(snapshot_time) 
      FROM issues_snapshot 
      WHERE repo_owner = :owner AND repo_name = :repo
  );

读取策略3:增量清洗(针对未清洗的数据)

SELECT s.*
FROM issues_snapshot s
LEFT JOIN experimental_issues e 
    ON s.issue_id = e.issue_id
WHERE s.repo_owner = :owner 
  AND s.repo_name = :repo
  AND e.issue_id IS NULL  -- 未清洗过的
  AND s.snapshot_time = :latest_time;

输出

  • raw_issues (List[Dict]): 原始Issue数据列表
  • total_count (Integer): 读取的Issue总数

设计代码结构

class DataCleaner:
    def load_issues(
        self, 
        repo_owner: str, 
        repo_name: str,
        start_date: Optional[date] = None,
        end_date: Optional[date] = None
    ) -> List[Dict]:
        """加载需要清洗的Issue"""
        # 实现逻辑
        pass

步骤2: 数据清洗和规范化

步骤ID: WF-005-S02
功能: 去重、格式规范化、缺失值处理
实现状态: ❌ 待实现

输入

  • raw_issues (步骤1输出): 原始Issue数据
  • cleaning_rules: 清洗规则配置

处理逻辑

2.1 去重处理

def deduplicate_issues(issues: List[Dict], method: str = "issue_id") -> List[Dict]:
    """去重:按issue_id保留最新记录"""
    seen = {}
    for issue in sorted(issues, key=lambda x: x.get('updated_at', '')):
        issue_id = issue.get('issue_id')
        if issue_id:
            seen[issue_id] = issue
    
    duplicates_removed = len(issues) - len(seen)
    print(f"✓ 去重: 移除 {duplicates_removed} 条重复记录")
    
    return list(seen.values())

2.2 字段规范化

def normalize_fields(issue: Dict, rules: Dict) -> Dict:
    """字段格式规范化"""
    cleaned = issue.copy()
    
    # 标题规范化
    if 'title' in cleaned:
        title = cleaned['title']
        # 去除首尾空格
        title = title.strip()
        # 限制长度
        max_len = rules.get('title', {}).get('max_length', 256)
        title = title[:max_len]
        cleaned['title'] = title
    
    # 状态规范化
    if 'state' in cleaned:
        state = cleaned['state'].lower()
        valid_states = rules.get('state', {}).get('valid_values', ['open', 'closed'])
        if state not in valid_states:
            cleaned['state'] = rules.get('state', {}).get('default', 'open')
    
    # Labels格式化(确保是JSON数组)
    if 'labels' in cleaned:
        labels = cleaned['labels']
        if isinstance(labels, str):
            try:
                cleaned['labels'] = json.loads(labels)
            except:
                cleaned['labels'] = []
        elif not isinstance(labels, list):
            cleaned['labels'] = []
    
    return cleaned

2.3 缺失值处理

def handle_missing_values(issue: Dict, rules: Dict) -> Dict:
    """处理缺失值"""
    cleaned = issue.copy()
    
    for field, config in rules.get('missing_values', {}).items():
        if field not in cleaned or cleaned[field] is None:
            action = config.get('action', 'keep_null')
            
            if action == 'keep_null':
                # 保持NULL
                pass
            
            elif action == 'set_default':
                # 设置默认值
                cleaned[field] = config.get('default')
            
            elif action == 'ai_infer':
                # 标记为需要AI推断(步骤3处理)
                cleaned[f'_{field}_needs_ai'] = True
    
    return cleaned

2.4 数据补全

def supplement_data(issue: Dict) -> Dict:
    """补全可以推断的数据"""
    # 如果没有issue_number,从其他字段推断
    if not issue.get('issue_number') and issue.get('issue_url'):
        match = re.search(r'/issues/(\d+)', issue['issue_url'])
        if match:
            issue['issue_number'] = int(match.group(1))
    
    # 如果没有repo_owner/repo_name,从URL推断
    if not issue.get('repo_owner') and issue.get('issue_url'):
        match = re.search(r'github\.com/([^/]+)/([^/]+)/issues', issue['issue_url'])
        if match:
            issue['repo_owner'] = match.group(1)
            issue['repo_name'] = match.group(2)
    
    return issue

输出

  • cleaned_issues (List[Dict]): 清洗后的Issue数据
  • cleaning_stats (Dict): 清洗统计信息

设计代码结构

def clean_data(
    self, 
    raw_issues: List[Dict],
    rules: Dict
) -> Tuple[List[Dict], Dict]:
    """数据清洗和规范化"""
    # 实现逻辑
    pass

步骤3: AI重新打标签

步骤ID: WF-005-S03
功能: 使用AI重新分析Issue,规范化类型、优先级、Labels
实现状态: ❌ 待实现

输入

  • cleaned_issues (步骤2输出): 清洗后的Issue
  • label_mapping: Labels标准化映射
  • AI配置: 通义千问API配置

处理逻辑

3.1 Labels标准化映射

def standardize_labels(issue: Dict, mapping: Dict) -> Dict:
    """应用Labels标准化映射"""
    labels = issue.get('labels', [])
    standardized = []
    
    for label in labels:
        label_name = label.get('name', label) if isinstance(label, dict) else label
        
        # 应用映射
        if label_name in mapping:
            standardized.append(mapping[label_name])
        # 跳过废弃标签
        elif label_name not in mapping.get('deprecated', []):
            standardized.append(label_name)
    
    issue['labels'] = standardized
    return issue

3.2 AI批量重新分析

def ai_relabel_batch(
    self, 
    issues: List[Dict], 
    batch_size: int = 50
) -> List[Dict]:
    """AI批量重新打标签"""
    results = []
    
    for i in range(0, len(issues), batch_size):
        batch = issues[i:i+batch_size]
        print(f"处理批次 {i//batch_size + 1}/{(len(issues)-1)//batch_size + 1}")
        
        # 为每个Issue调用AI
        for issue in batch:
            try:
                # 构建Prompt
                prompt = self._build_relabel_prompt(issue)
                
                # 调用AI
                ai_response = self.llm._call_ai(
                    system_prompt="你是Issue标注专家,负责规范化Issue的分类和标签。",
                    user_prompt=prompt
                )
                
                # 解析AI响应
                ai_labels = self._parse_ai_response(ai_response)
                
                # 更新Issue
                issue['ai_issue_type'] = ai_labels.get('issue_type')
                issue['ai_priority'] = ai_labels.get('priority')
                issue['ai_labels'] = ai_labels.get('labels', [])
                issue['ai_corrected'] = True
                
                results.append(issue)
                
            except Exception as e:
                print(f"⚠️ Issue #{issue.get('issue_number')} AI分析失败: {e}")
                issue['ai_corrected'] = False
                results.append(issue)
        
        # 避免API限流
        time.sleep(1)
    
    return results

3.3 AI Prompt构建

def _build_relabel_prompt(self, issue: Dict) -> str:
    """构建AI重新标注的Prompt"""
    return f"""
请重新分析以下Issue,规范化其分类和标签。

【Issue信息】
标题: {issue.get('title')}
正文: {issue.get('body', '')[:500]}...
当前Labels: {issue.get('labels', [])}

【任务】
1. 判断Issue类型(bug/feature/task/question)
2. 判断优先级(P0/P1/P2/P3)
3. 规范化Labels(使用标准前缀:kind/, area/, severity/)

【规则】
- Labels必须使用标准前缀
- 优先级基于影响范围和紧急程度
- 如果标题/正文含客户信息,添加customer/标签

请返回JSON格式:
{{
    "issue_type": "bug",
    "priority": "P1",
    "labels": ["kind/bug", "area/问数", "severity/high"]
}}
"""

输出

  • relabeled_issues (List[Dict]): AI重新标注的Issue
  • ai_stats (Dict): AI处理统计
    {
        "total_processed": 4950,
        "success_count": 4900,
        "failed_count": 50,
        "labels_modified": 1500,
        "average_processing_time": 2.5
    }

设计代码结构

def ai_relabel(
    self, 
    issues: List[Dict],
    mapping: Dict
) -> Tuple[List[Dict], Dict]:
    """AI重新打标签"""
    # 实现逻辑
    pass

步骤4: 质量验证

步骤ID: WF-005-S04
功能: 检查清洗后数据的完整性和质量
实现状态: ❌ 待实现

输入

  • relabeled_issues (步骤3输出): 重新标注的Issue
  • quality_rules: 质量检查规则

处理逻辑

4.1 必需字段检查

def validate_required_fields(issue: Dict, required: List[str]) -> Tuple[bool, List[str]]:
    """检查必需字段"""
    missing = []
    for field in required:
        if field not in issue or issue[field] is None or issue[field] == '':
            missing.append(field)
    
    is_valid = len(missing) == 0
    return is_valid, missing

4.2 Labels有效性检查

def validate_labels(issue: Dict, valid_prefixes: List[str]) -> Tuple[bool, List[str]]:
    """检查Labels格式"""
    labels = issue.get('labels', [])
    invalid = []
    
    for label in labels:
        # 检查是否有有效前缀
        has_valid_prefix = any(label.startswith(prefix) for prefix in valid_prefixes)
        if not has_valid_prefix and '/' in label:
            invalid.append(label)
    
    is_valid = len(invalid) == 0
    return is_valid, invalid

4.3 数据质量评分

def calculate_quality_score(issue: Dict) -> float:
    """计算数据质量分数(0-1)"""
    score = 0.0
    max_score = 0.0
    
    # 标题质量(20分)
    max_score += 20
    if issue.get('title'):
        title_len = len(issue['title'])
        if 10 <= title_len <= 200:
            score += 20
        elif title_len > 0:
            score += 10
    
    # Labels质量(30分)
    max_score += 30
    labels = issue.get('labels', [])
    if len(labels) >= 2:  # 至少2个标签
        score += 15
    has_kind = any('kind/' in l for l in labels)
    has_area = any('area/' in l for l in labels)
    if has_kind:
        score += 10
    if has_area:
        score += 5
    
    # AI分析质量(25分)
    max_score += 25
    if issue.get('ai_corrected'):
        score += 15
    if issue.get('ai_issue_type'):
        score += 5
    if issue.get('ai_priority'):
        score += 5
    
    # 正文质量(15分)
    max_score += 15
    if issue.get('body'):
        body_len = len(issue.get('body', ''))
        if body_len >= 50:
            score += 15
        elif body_len > 0:
            score += 5
    
    # 其他字段完整性(10分)
    max_score += 10
    if issue.get('assignee'):
        score += 5
    if issue.get('milestone'):
        score += 5
    
    return score / max_score if max_score > 0 else 0.0

4.4 生成质量报告

def generate_quality_report(issues: List[Dict], min_score: float) -> Dict:
    """生成质量报告"""
    quality_scores = [calculate_quality_score(issue) for issue in issues]
    
    passed = [s for s in quality_scores if s >= min_score]
    failed = [s for s in quality_scores if s < min_score]
    
    report = {
        "total_issues": len(issues),
        "average_quality_score": sum(quality_scores) / len(quality_scores),
        "median_quality_score": sorted(quality_scores)[len(quality_scores) // 2],
        "min_quality_score": min(quality_scores),
        "max_quality_score": max(quality_scores),
        "passed_threshold": len(passed),
        "failed_threshold": len(failed),
        "pass_rate": len(passed) / len(issues)
    }
    
    return report

输出

  • validated_issues (List[Dict]): 验证通过的Issue
  • failed_issues (List[Dict]): 验证失败的Issue
  • quality_report (Dict): 质量报告

设计代码结构

def validate_quality(
    self, 
    issues: List[Dict],
    rules: Dict
) -> Tuple[List[Dict], List[Dict], Dict]:
    """质量验证"""
    # 实现逻辑
    pass

步骤5: 存入实验库

步骤ID: WF-005-S05
功能: 将清洗后的数据存入experimental_issues表
实现状态: ❌ 待实现

输入

  • validated_issues (步骤4输出): 验证通过的Issue
  • cleaning_metadata: 清洗元数据

处理逻辑

5.1 创建实验库表

CREATE TABLE IF NOT EXISTS experimental_issues (
    id INTEGER PRIMARY KEY AUTO_INCREMENT,
    issue_id BIGINT NOT NULL,
    issue_number INTEGER NOT NULL,
    repo_owner VARCHAR(100) NOT NULL,
    repo_name VARCHAR(100) NOT NULL,
    
    -- 清洗后的字段
    title VARCHAR(256) NOT NULL,
    body TEXT,
    state VARCHAR(20) NOT NULL,
    labels JSON,
    assignee VARCHAR(100),
    milestone VARCHAR(100),
    
    -- AI分析结果
    ai_issue_type VARCHAR(50),
    ai_priority VARCHAR(10),
    ai_labels JSON,
    ai_corrected BOOLEAN DEFAULT FALSE,
    
    -- 数据质量
    quality_score FLOAT,
    validation_passed BOOLEAN DEFAULT TRUE,
    
    -- 清洗元数据
    cleaned_at DATETIME NOT NULL,
    cleaning_version VARCHAR(50),
    source_snapshot_time DATETIME,
    
    -- GitHub时间戳
    created_at DATETIME,
    updated_at DATETIME,
    closed_at DATETIME,
    
    -- 索引
    INDEX idx_issue_id (issue_id),
    INDEX idx_repo (repo_owner, repo_name),
    INDEX idx_cleaned_at (cleaned_at),
    UNIQUE KEY uk_issue_cleaning (issue_id, cleaning_version)
);

5.2 批量插入

def save_to_experimental(
    self, 
    issues: List[Dict],
    cleaning_version: str
) -> Dict:
    """保存到实验库"""
    success_count = 0
    error_count = 0
    
    for issue in issues:
        try:
            sql = """
            INSERT INTO experimental_issues (
                issue_id, issue_number, repo_owner, repo_name,
                title, body, state, labels, assignee, milestone,
                ai_issue_type, ai_priority, ai_labels, ai_corrected,
                quality_score, validation_passed,
                cleaned_at, cleaning_version, source_snapshot_time,
                created_at, updated_at, closed_at
            ) VALUES (
                :issue_id, :issue_number, :repo_owner, :repo_name,
                :title, :body, :state, :labels, :assignee, :milestone,
                :ai_issue_type, :ai_priority, :ai_labels, :ai_corrected,
                :quality_score, :validation_passed,
                :cleaned_at, :cleaning_version, :source_snapshot_time,
                :created_at, :updated_at, :closed_at
            )
            ON DUPLICATE KEY UPDATE
                title = VALUES(title),
                body = VALUES(body),
                state = VALUES(state),
                labels = VALUES(labels),
                ai_issue_type = VALUES(ai_issue_type),
                ai_priority = VALUES(ai_priority),
                ai_labels = VALUES(ai_labels),
                quality_score = VALUES(quality_score),
                cleaned_at = VALUES(cleaned_at)
            """
            
            self.storage.execute(sql, {
                "issue_id": issue['issue_id'],
                "issue_number": issue['issue_number'],
                "repo_owner": issue['repo_owner'],
                "repo_name": issue['repo_name'],
                "title": issue['title'],
                "body": issue.get('body'),
                "state": issue['state'],
                "labels": json.dumps(issue.get('labels', [])),
                "assignee": issue.get('assignee'),
                "milestone": issue.get('milestone'),
                "ai_issue_type": issue.get('ai_issue_type'),
                "ai_priority": issue.get('ai_priority'),
                "ai_labels": json.dumps(issue.get('ai_labels', [])),
                "ai_corrected": issue.get('ai_corrected', False),
                "quality_score": issue.get('quality_score', 0.0),
                "validation_passed": issue.get('validation_passed', True),
                "cleaned_at": datetime.now(),
                "cleaning_version": cleaning_version,
                "source_snapshot_time": issue.get('snapshot_time'),
                "created_at": issue.get('created_at'),
                "updated_at": issue.get('updated_at'),
                "closed_at": issue.get('closed_at')
            })
            
            success_count += 1
            
        except Exception as e:
            print(f"⚠️ 保存Issue #{issue.get('issue_number')} 失败: {e}")
            error_count += 1
    
    return {
        "success_count": success_count,
        "error_count": error_count,
        "total": len(issues)
    }

5.3 记录变更日志

# 创建变更日志表
CREATE TABLE IF NOT EXISTS data_cleaning_log (
    id INTEGER PRIMARY KEY AUTO_INCREMENT,
    issue_id BIGINT NOT NULL,
    cleaning_version VARCHAR(50) NOT NULL,
    action_type VARCHAR(50),  -- 'normalize', 'ai_relabel', 'quality_check'
    field_name VARCHAR(100),
    old_value TEXT,
    new_value TEXT,
    created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
    INDEX idx_issue_id (issue_id),
    INDEX idx_cleaning_version (cleaning_version)
);

# 插入变更记录
def log_change(
    issue_id: int,
    cleaning_version: str,
    action: str,
    field: str,
    old_value: Any,
    new_value: Any
):
    """记录数据变更"""
    sql = """
    INSERT INTO data_cleaning_log 
        (issue_id, cleaning_version, action_type, field_name, old_value, new_value)
    VALUES 
        (:issue_id, :version, :action, :field, :old, :new)
    """
    self.storage.execute(sql, {
        "issue_id": issue_id,
        "version": cleaning_version,
        "action": action,
        "field": field,
        "old": str(old_value)[:1000],
        "new": str(new_value)[:1000]
    })

5.4 生成清洗报告

def generate_cleaning_report(
    self,
    stats: Dict,
    quality_report: Dict
) -> str:
    """生成Markdown格式的清洗报告"""
    report = f"""
# 数据清洗报告

## 执行信息
- 开始时间: {stats['start_time']}
- 结束时间: {stats['end_time']}
- 总耗时: {stats['duration_seconds']}
- 仓库: {stats['repo_owner']}/{stats['repo_name']}

## 数据统计
- 读取Issue总数: {stats['total_issues_read']}
- 去重移除: {stats['duplicates_removed']}
- 成功清洗: {stats['records_cleaned']}
- 清洗失败: {stats['records_failed']}
- 成功率: {stats['success_rate']:.2%}

## 清洗操作
- 字段规范化: {stats['field_normalization']}
- Labels标准化: {stats['label_standardization']}
- AI重新标注: {stats['ai_re_labeling']}
- 缺失值填充: {stats['missing_value_filled']}

## 数据质量
- 平均质量分: {quality_report['average_quality_score']:.2f}
- 中位数分数: {quality_report['median_quality_score']:.2f}
- 通过阈值: {quality_report['passed_threshold']} ({quality_report['pass_rate']:.2%})
- 未达标: {quality_report['failed_threshold']}

## Labels变更
- 总修改数: {stats['total_labels_modified']}
- 映射应用: {stats['mapping_applied']}
- AI校正: {stats['ai_corrections']}
"""
    
    return report

输出

  • 数据库记录: 成功写入experimental_issues表
  • cleaning_report.md: Markdown格式报告
  • cleaning_report.json: JSON格式报告
  • save_stats: 保存统计信息

设计代码结构

def save_and_report(
    self, 
    issues: List[Dict],
    cleaning_stats: Dict,
    quality_report: Dict
) -> Dict:
    """保存数据并生成报告"""
    # 实现逻辑
    pass

🗄️ 数据库表结构

experimental_issues(实验Issue表)

说明: 存储清洗和重新标注后的Issue数据

字段名 类型 说明 索引
id INTEGER 主键 PK
issue_id BIGINT GitHub Issue ID YES
issue_number INTEGER Issue编号 -
repo_owner VARCHAR(100) 仓库所有者 YES
repo_name VARCHAR(100) 仓库名称 YES
title VARCHAR(256) 清洗后标题 -
body TEXT 清洗后正文 -
state VARCHAR(20) 状态 -
labels JSON 标准化Labels -
assignee VARCHAR(100) 负责人 -
milestone VARCHAR(100) 里程碑 -
ai_issue_type VARCHAR(50) AI判断的类型 -
ai_priority VARCHAR(10) AI判断的优先级 -
ai_labels JSON AI推荐的Labels -
ai_corrected BOOLEAN 是否经过AI校正 -
quality_score FLOAT 数据质量分数(0-1) -
validation_passed BOOLEAN 是否通过验证 -
cleaned_at DATETIME 清洗时间 YES
cleaning_version VARCHAR(50) 清洗版本号 -
source_snapshot_time DATETIME 源快照时间 -
created_at DATETIME GitHub创建时间 -
updated_at DATETIME GitHub更新时间 -
closed_at DATETIME GitHub关闭时间 -

唯一约束: (issue_id, cleaning_version)


data_cleaning_log(数据清洗日志表)

说明: 记录清洗过程中的所有数据变更

字段名 类型 说明 索引
id INTEGER 主键 PK
issue_id BIGINT Issue ID YES
cleaning_version VARCHAR(50) 清洗版本 YES
action_type VARCHAR(50) 操作类型 -
field_name VARCHAR(100) 字段名 -
old_value TEXT 旧值 -
new_value TEXT 新值 -
created_at DATETIME 记录时间 -

操作类型

  • normalize: 字段规范化
  • ai_relabel: AI重新标注
  • quality_check: 质量检查
  • mapping: Labels映射

⚙️ 配置文件说明

配置文件位置

主配置: config/cleaning_rules.yaml(新建)

配置示例已在"整体输入"部分展示。


🔧 运行方式(设计)

方式1:命令行运行

python3 scripts/clean_historical_data.py \
    --repo-owner matrixorigin \
    --repo-name matrixone \
    --start-date 2024-01-01 \
    --end-date 2024-12-31 \
    --config config/cleaning_rules.yaml

方式2:Python代码调用

from modules.data_cleaning.cleaner import DataCleaner

# 初始化
cleaner = DataCleaner(
    storage=MOStorage(),
    llm=LLMParser(),
    config_path='config/cleaning_rules.yaml'
)

# 执行清洗
result = cleaner.run(
    repo_owner='matrixorigin',
    repo_name='matrixone',
    start_date=date(2024, 1, 1),
    end_date=date(2024, 12, 31)
)

print(f"✅ 清洗完成: {result['records_cleaned']} 条记录")
print(f"质量分数: {result['average_quality_score']:.2f}")

✅ 实现建议

开发优先级

P0(核心功能)

  1. 步骤1: 数据读取
  2. 步骤2: 字段规范化和去重
  3. 步骤5: 保存到实验库

P1(AI增强): 4. 步骤3: AI重新打标签 5. Labels标准化映射

P2(质量保证): 6. 步骤4: 质量验证 7. 清洗报告生成 8. 变更日志记录

技术选型

  • 数据处理: pandas(批量处理)
  • AI调用: 复用 modules/llm_parser/llm_parser.py
  • 配置管理: PyYAML
  • 数据库: 复用 modules/database_storage/mo_client.py

性能优化

  • 批量处理: 每批50-100条Issue
  • 并行AI调用: 使用线程池(注意API限流)
  • 增量清洗: 只处理未清洗的数据
  • 缓存: 缓存AI响应避免重复调用

🔄 云端部署注意事项

✅ 完全可用

该工作流设计上在云端部署无任何限制

  • ✅ 纯数据库+AI API操作
  • ✅ 无本地文件依赖
  • ✅ 可配置化

部署建议

  1. 分批执行:避免一次性处理过多数据
  2. 定时任务:可设置为每周执行一次
  3. 监控告警:监控清洗失败率和质量分数

📊 预期效果

指标 预期值 说明
数据完整性 >95% 必需字段齐全
Labels规范率 >90% 使用标准前缀
AI校正准确率 >85% AI重新标注准确
平均质量分 >0.85 整体数据质量
处理速度 ~100条/分钟 包含AI调用

文档版本: v1.0
最后更新: 2026-03-04
状态: 设计文档(待实现)