把数据质量换到数据安全时踩过的坑

把数据质量换到数据安全时踩过的坑上手并不难,难的是稳定跑起来。

下面只记真正影响结果的部分。

数据质量

数据完整性

-- 检查缺失值
SELECT 
  COUNT(*) as total,
  SUM(CASE WHEN email IS NULL THEN 1 ELSE 0 END) as missing_email,
  SUM(CASE WHEN phone IS NULL THEN 1 ELSE 0 END) as missing_phone
FROM users;

-- 计算完整度
SELECT 
  100 * (
    SUM(CASE WHEN email IS NOT NULL AND phone IS NOT NULL THEN 1 ELSE 0 END) / 
    COUNT(*)
  ) as completeness_rate
FROM users;

数据准确性

-- 检查异常值
SELECT 
  age,
  COUNT(*) as count
FROM users
GROUP BY age
HAVING age < 0 OR age > 120;

-- 检查重复数据
SELECT 
  email,
  COUNT(*) as count
FROM users
GROUP BY email
HAVING COUNT(*) > 1;

-- 检查格式错误
SELECT 
  email
FROM users
WHERE email NOT REGEXP '^[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\\.[A-Z|a-z]{2,}$';

数据一致性

-- 检查外键一致性
SELECT 
  o.user_id,
  u.id as user_exists
FROM orders o
LEFT JOIN users u ON o.user_id = u.id
WHERE u.id IS NULL;

-- 检查数据状态一致性
SELECT 
  status,
  COUNT(*) as count
FROM orders
WHERE status NOT IN ('pending', 'processing', 'completed', 'cancelled')
GROUP BY status;

数据清洗

数据清洗流程

import pandas as pd
import numpy as np

def clean_data(df):
    # 1. 处理缺失值
    df = handle_missing_values(df)
    
    # 2. 处理异常值
    df = handle_outliers(df)
    
    # 3. 数据标准化
    df = standardize_data(df)
    
    # 4. 数据去重
    df = remove_duplicates(df)
    
    return df

def handle_missing_values(df):
    # 数值列:填充中位数
    numeric_cols = df.select_dtypes(include=[np.number]).columns
    df[numeric_cols] = df[numeric_cols].fillna(df[numeric_cols].median())
    
    # 分类列:填充众数
    categorical_cols = df.select_dtypes(include=['object']).columns
    for col in categorical_cols:
        df[col].fillna(df[col].mode()[0], inplace=True)
    
    return df

def handle_outliers(df):
    # 使用 IQR 方法检测异常值
    numeric_cols = df.select_dtypes(include=[np.number]).columns
    
    for col in numeric_cols:
        Q1 = df[col].quantile(0.25)
        Q3 = df[col].quantile(0.75)
        IQR = Q3 - Q1
        
        lower_bound = Q1 - 1.5 * IQR
        upper_bound = Q3 + 1.5 * IQR
        
        # 用边界值替换异常值
        df[col] = np.where(df[col] < lower_bound, lower_bound, df[col])
        df[col] = np.where(df[col] > upper_bound, upper_bound, df[col])
    
    return df

def standardize_data(df):
    # 标准化格式
    df['email'] = df['email'].str.lower()
    df['phone'] = df['phone'].str.replace(r'[^0-9]', '')
    
    return df

def remove_duplicates(df):
    # 去除完全重复的行
    df = df.drop_duplicates()
    
    # 去除关键字段重复的行
    df = df.drop_duplicates(subset=['email'])
    
    return df

自动化数据质量检查

from typing import List, Dict, Any

class DataQualityChecker:
    def __init__(self, rules: List[Dict[str, Any]]):
        self.rules = rules
    
    def check(self, df: pd.DataFrame) -> Dict[str, Any]:
        results = {}
        
        for rule in self.rules:
            rule_name = rule['name']
            rule_type = rule['type']
            
            if rule_type == 'completeness':
                results[rule_name] = self.check_completeness(df, rule)
            elif rule_type == 'uniqueness':
                results[rule_name] = self.check_uniqueness(df, rule)
            elif rule_type == 'validity':
                results[rule_name] = self.check_validity(df, rule)
            elif rule_type == 'consistency':
                results[rule_name] = self.check_consistency(df, rule)
        
        return results
    
    def check_completeness(self, df: pd.DataFrame, rule: Dict[str, Any]) -> Dict[str, Any]:
        column = rule['column']
        threshold = rule['threshold']
        
        completeness = df[column].notna().sum() / len(df)
        passed = completeness >= threshold
        
        return {
            'passed': passed,
            'actual': completeness,
            'expected': threshold
        }
    
    def check_uniqueness(self, df: pd.DataFrame, rule: Dict[str, Any]) -> Dict[str, Any]:
        column = rule['column']
        threshold = rule['threshold']
        
        uniqueness = df[column].nunique() / len(df)
        passed = uniqueness >= threshold
        
        return {
            'passed': passed,
            'actual': uniqueness,
            'expected': threshold
        }

# 使用规则
rules = [
    {
        'name': 'email_completeness',
        'type': 'completeness',
        'column': 'email',
        'threshold': 0.95
    },
    {
        'name': 'email_uniqueness',
        'type': 'uniqueness',
        'column': 'email',
        'threshold': 0.99
    }
]

checker = DataQualityChecker(rules)
results = checker.check(df)

数据安全

数据加密

from cryptography.fernet import Fernet

class DataEncryptor:
    def __init__(self, key: bytes):
        self.cipher = Fernet(key)
    
    def encrypt(self, data: str) -> bytes:
        return self.cipher.encrypt(data.encode())
    
    def decrypt(self, encrypted_data: bytes) -> str:
        return self.cipher.decrypt(encrypted_data).decode()

# 使用
key = Fernet.generate_key()
encryptor = DataEncryptor(key)

encrypted = encryptor.encrypt("sensitive data")
decrypted = encryptor.decrypt(encrypted)

数据脱敏

import re

class DataMasker:
    @staticmethod
    def mask_email(email: str) -> str:
        if not email:
            return email
        
        username, domain = email.split('@')
        masked_username = username[:2] + '*' * (len(username) - 2)
        
        return f"{masked_username}@{domain}"
    
    @staticmethod
    def mask_phone(phone: str) -> str:
        if not phone:
            return phone
        
        return phone[:3] + '*' * (len(phone) - 6) + phone[-3:]
    
    @staticmethod
    def mask_card_number(card_number: str) -> str:
        if not card_number:
            return card_number
        
        return '*' * (len(card_number) - 4) + card_number[-4:]

# 使用
masker = DataMasker()
print(masker.mask_email("[email protected]"))  # al***@example.com
print(masker.mask_phone("13812345678"))  # 138****5678
print(masker.mask_card_number("1234567890123456"))  # ************3456

访问控制

from functools import wraps
from typing import Callable

class AccessControl:
    def __init__(self):
        self.permissions = {}
    
    def add_permission(self, role: str, resource: str, action: str):
        if role not in self.permissions:
            self.permissions[role] = []
        
        self.permissions[role].append(f"{resource}:{action}")
    
    def check_permission(self, role: str, resource: str, action: str) -> bool:
        if role not in self.permissions:
            return False
        
        return f"{resource}:{action}" in self.permissions[role]
    
    def require_permission(self, resource: str, action: str):
        def decorator(func: Callable):
            @wraps(func)
            def wrapper(*args, **kwargs):
                # 假设从上下文中获取用户角色
                role = get_user_role()
                
                if not self.check_permission(role, resource, action):
                    raise PermissionError("Access denied")
                
                return func(*args, **kwargs)
            return wrapper
        return decorator

# 使用
access_control = AccessControl()
access_control.add_permission('admin', 'users', 'read')
access_control.add_permission('admin', 'users', 'write')

@access_control.require_permission('users', 'read')
def get_users():
    return db.query('SELECT * FROM users')

数据血缘

追踪数据来源

class DataLineage:
    def __init__(self):
        self.lineage = {}
    
    def add_source(self, data_id: str, source_id: str, metadata: dict = None):
        if data_id not in self.lineage:
            self.lineage[data_id] = []
        
        self.lineage[data_id].append({
            'source_id': source_id,
            'metadata': metadata or {},
            'timestamp': datetime.now()
        })
    
    def get_lineage(self, data_id: str) -> list:
        return self.lineage.get(data_id, [])
    
    def trace_back(self, data_id: str, max_depth: int = 10) -> list:
        lineage = []
        visited = set()
        
        self._trace_back(data_id, lineage, visited, 0, max_depth)
        
        return lineage
    
    def _trace_back(self, data_id: str, lineage: list, visited: set, depth: int, max_depth: int):
        if data_id in visited or depth >= max_depth:
            return
        
        visited.add(data_id)
        
        sources = self.lineage.get(data_id, [])
        for source in sources:
            lineage.append({
                'data_id': data_id,
                'source_id': source['source_id'],
                'metadata': source['metadata'],
                'depth': depth
            })
            
            self._trace_back(source['source_id'], lineage, visited, depth + 1, max_depth)

# 使用
lineage = DataLineage()
lineage.add_source('report_1', 'data_source_1', {'timestamp': '2023-01-01'})
lineage.add_source('report_1', 'data_source_2', {'timestamp': '2023-01-02'})

print(lineage.get_lineage('report_1'))
print(lineage.trace_back('report_1'))

数据治理工具

数据治理框架

from typing import Dict, Any, List

class DataGovernanceFramework:
    def __init__(self):
        self.policies = []
        self.audit_log = []
    
    def add_policy(self, policy: Dict[str, Any]):
        self.policies.append(policy)
    
    def audit(self, df: pd.DataFrame) -> Dict[str, Any]:
        results = {
            'passed': True,
            'violations': []
        }
        
        for policy in self.policies:
            try:
                if policy['type'] == 'data_quality':
                    violations = self._check_data_quality(df, policy)
                    if violations:
                        results['passed'] = False
                        results['violations'].extend(violations)
                
                elif policy['type'] == 'data_security':
                    violations = self._check_data_security(df, policy)
                    if violations:
                        results['passed'] = False
                        results['violations'].extend(violations)
            
            except Exception as e:
                results['passed'] = False
                results['violations'].append({
                    'policy': policy['name'],
                    'error': str(e)
                })
        
        self.audit_log.append({
            'timestamp': datetime.now(),
            'results': results
        })
        
        return results
    
    def _check_data_quality(self, df: pd.DataFrame, policy: Dict[str, Any]) -> List[Dict[str, Any]]:
        violations = []
        rule = policy['rule']
        
        if rule == 'no_missing_values':
            missing = df.isna().sum()
            for col, count in missing.items():
                if count > 0:
                    violations.append({
                        'policy': policy['name'],
                        'type': 'missing_values',
                        'column': col,
                        'count': count
                    })
        
        elif rule == 'no_duplicates':
            duplicates = df.duplicated().sum()
            if duplicates > 0:
                violations.append({
                    'policy': policy['name'],
                    'type': 'duplicates',
                    'count': duplicates
                })
        
        return violations
    
    def get_audit_report(self) -> Dict[str, Any]:
        passed = sum(1 for log in self.audit_log if log['results']['passed'])
        total = len(self.audit_log)
        
        return {
            'total_audits': total,
            'passed_audits': passed,
            'failed_audits': total - passed,
            'pass_rate': passed / total if total > 0 else 0,
            'recent_audits': self.audit_log[-10:]
        }

# 使用
framework = DataGovernanceFramework()
framework.add_policy({
    'name': 'no_missing_email',
    'type': 'data_quality',
    'rule': 'no_missing_values'
})

results = framework.audit(df)
report = framework.get_audit_report()

踩过的坑

坑一:数据清洗过度

数据清洗过度,导致有用数据丢失。

解决:记录清洗过程,可以回滚。

class DataCleaner:
    def __init__(self):
        self.cleaning_history = []
    
    def clean(self, df: pd.DataFrame) -> pd.DataFrame:
        original = df.copy()
        
        # 清洗数据
        df = self.handle_missing_values(df)
        df = self.remove_duplicates(df)
        
        # 记录清洗历史
        self.cleaning_history.append({
            'timestamp': datetime.now(),
            'original_rows': len(original),
            'cleaned_rows': len(df),
            'removed_rows': len(original) - len(df)
        })
        
        return df
    
    def rollback(self, df: pd.DataFrame) -> pd.DataFrame:
        # 可以实现回滚逻辑
        pass

坑二:数据脱敏不彻底

敏感信息没有被完全脱敏。

解决:多层脱敏,定期检查。

class DataMasker:
    def mask(self, data: dict, schema: dict) -> dict:
        masked = data.copy()
        
        for field, mask_type in schema.items():
            if field in masked:
                if mask_type == 'email':
                    masked[field] = self.mask_email(masked[field])
                elif mask_type == 'phone':
                    masked[field] = self.mask_phone(masked[field])
                elif mask_type == 'hash':
                    masked[field] = hashlib.sha256(masked[field].encode()).hexdigest()
        
        return masked

# 使用
schema = {
    'email': 'email',
    'phone': 'phone',
    'ssn': 'hash'
}

masked_data = masker.mask(user_data, schema)

坑三:权限管理混乱

权限设置不合理,导致数据泄露。

解决:最小权限原则,定期审计。

class PermissionManager:
    def audit_permissions(self) -> list:
        violations = []
        
        for role in self.roles:
            permissions = self.get_role_permissions(role)
            
            # 检查是否有不必要的权限
            unnecessary = self.check_unnecessary_permissions(permissions)
            if unnecessary:
                violations.append({
                    'role': role,
                    'unnecessary_permissions': unnecessary
                })
        
        return violations

写在最后

数据治理这东西,不只是技术问题,是管理和文化问题。

解决了

  • 数据质量问题
  • 数据安全问题
  • 数据管理问题

带来了

  • 管理成本
  • 技术成本
  • 复杂度增加

实施数据治理之前先评估:

  • 数据规模
  • 业务需求
  • 团队能力
  • 法规要求

不是所有数据都需要严格治理,但关键数据不能马虎。


这次数据治理花了一个月,从数据质量到数据安全。治理完成后,数据质量提升了 80%,安全事件减少了 90%。

版权声明: 本文首发于 指尖魔法屋-把数据质量换到数据安全时踩过的坑https://blog.thinkmoon.cn/post/82-data-governance-quality-security-practice/) 转载或引用必须申明原指尖魔法屋来源及源地址!