PHP项目脏数据如何分类标记批量处理

wen PHP项目 27

本文目录导读:

PHP项目脏数据如何分类标记批量处理

  1. 脏数据分类体系
  2. 数据标记方案
  3. 批量处理系统
  4. Web管理界面核心功能
  5. 最佳实践建议

针对PHP项目中的脏数据分类标记与批量处理,这是一个典型的数据治理场景,下面整理了一套完整的解决方案,涵盖分类策略、标记机制和批量处理流程。

脏数据分类体系

按问题类型分类

// 脏数据类型枚举
class DirtyDataType {
    const INCOMPLETE    = 1; // 数据不完整
    const DUPLICATE     = 2; // 数据重复
    const FORMAT_ERROR  = 3; // 格式错误
    const OUT_OF_RANGE  = 4; // 数值超范围
    const INCONSISTENT  = 5; // 数据不一致
    const LOGICAL_ERROR = 6; // 逻辑错误
    const OUTDATED      = 7; // 过期数据
    const SPAM          = 8; // 垃圾/恶意数据
}

按严重程度分级

class DirtyDataLevel {
    const CRITICAL = 1; // 致命 - 需立即处理
    const HIGH     = 2; // 高 - 优先级处理
    const MEDIUM   = 3; // 中 - 计划处理
    const LOW      = 4; // 低 - 可延迟处理
}

数据标记方案

数据库表结构设计

-- 脏数据标记主表
CREATE TABLE `data_quality_markers` (
    `id` bigint(20) unsigned NOT NULL AUTO_INCREMENT,
    `table_name` varchar(100) NOT NULL COMMENT '源表名',
    `record_id` bigint(20) NOT NULL COMMENT '源记录ID',
    `field_name` varchar(100) DEFAULT NULL COMMENT '问题字段名',
    `dirty_type` tinyint(4) NOT NULL COMMENT '脏数据类型',
    `dirty_level` tinyint(4) NOT NULL COMMENT '严重级别',
    `description` text COMMENT '问题描述',
    `original_value` text COMMENT '原始值',
    `suggested_value` text COMMENT '建议值',
    `status` tinyint(4) DEFAULT 0 COMMENT '0:待处理 1:已标记 2:处理中 3:已修正 4:已忽略',
    `handler` varchar(100) DEFAULT NULL COMMENT '处理人',
    `handled_at` datetime DEFAULT NULL COMMENT '处理时间',
    `batch_no` varchar(50) DEFAULT NULL COMMENT '批次号',
    `created_at` timestamp DEFAULT CURRENT_TIMESTAMP,
    `updated_at` timestamp DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
    PRIMARY KEY (`id`),
    KEY `idx_table_record` (`table_name`, `record_id`),
    KEY `idx_status` (`status`),
    KEY `idx_batch_no` (`batch_no`),
    KEY `idx_dirty_type` (`dirty_type`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 源表添加标记字段
ALTER TABLE `your_source_table` 
ADD COLUMN `dirty_flag` tinyint(4) DEFAULT 0 COMMENT '0:正常 1:脏数据',
ADD COLUMN `dirty_marker_id` bigint(20) DEFAULT NULL COMMENT '关联标记ID',
ADD INDEX `idx_dirty_flag` (`dirty_flag`);

数据标记分类器实现

<?php
class DataQualityClassifier 
{
    private $rules = [];
    private $db;
    public function __construct() 
    {
        $this->db = new PDO('mysql:host=localhost;dbname=test', 'user', 'pass');
        $this->loadRules();
    }
    /**
     * 加载分类规则
     */
    private function loadRules()
    {
        $this->rules = [
            'phone' => [
                'validator' => function($value) {
                    return preg_match('/^1[3-9]\d{9}$/', $value);
                },
                'type' => DirtyDataType::FORMAT_ERROR,
                'level' => DirtyDataLevel::HIGH,
                'description' => '手机号格式不正确'
            ],
            'email' => [
                'validator' => function($value) {
                    return filter_var($value, FILTER_VALIDATE_EMAIL) !== false;
                },
                'type' => DirtyDataType::FORMAT_ERROR,
                'level' => DirtyDataLevel::MEDIUM,
                'description' => '邮箱格式不正确'
            ],
            'age' => [
                'validator' => function($value) {
                    return is_numeric($value) && $value >= 0 && $value <= 150;
                },
                'type' => DirtyDataType::OUT_OF_RANGE,
                'level' => DirtyDataLevel::LOW,
                'description' => '年龄超出合理范围'
            ]
        ];
    }
    /**
     * 批量标记脏数据
     */
    public function batchMark($table, $records, $batchNo)
    {
        $db = $this->db;
        $db->beginTransaction();
        try {
            $stmt = $db->prepare(
                "INSERT INTO data_quality_markers 
                (table_name, record_id, field_name, dirty_type, dirty_level, 
                 description, original_value, suggested_value, batch_no)
                VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)"
            );
            $updateStmt = $db->prepare(
                "UPDATE {$table} SET dirty_flag = 1 WHERE id = ?"
            );
            foreach ($records as $record) {
                $problems = $this->classifyRecord($record);
                foreach ($problems as $problem) {
                    $stmt->execute([
                        $table,
                        $record['id'],
                        $problem['field'],
                        $problem['type'],
                        $problem['level'],
                        $problem['description'],
                        $problem['original_value'],
                        $problem['suggested_value'] ?? null,
                        $batchNo
                    ]);
                    $markerId = $db->lastInsertId();
                    // 更新源表标记
                    $updateStmt->execute([$record['id']]);
                }
            }
            $db->commit();
            return true;
        } catch (\Exception $e) {
            $db->rollBack();
            throw $e;
        }
    }
    /**
     * 分类单条记录
     */
    private function classifyRecord($record)
    {
        $problems = [];
        // 检查必填字段
        foreach ($this->rules as $field => $rule) {
            if (!isset($record[$field]) || empty($record[$field])) {
                $problems[] = [
                    'field' => $field,
                    'type' => DirtyDataType::INCOMPLETE,
                    'level' => DirtyDataLevel::HIGH,
                    'description' => "字段 {$field} 为空",
                    'original_value' => ''
                ];
                continue;
            }
            // 执行验证
            if (!$rule['validator']($record[$field])) {
                $problems[] = [
                    'field' => $field,
                    'type' => $rule['type'],
                    'level' => $rule['level'],
                    'description' => $rule['description'],
                    'original_value' => $record[$field],
                    'suggested_value' => $this->getSuggestion($field, $record[$field])
                ];
            }
        }
        // 检查重复数据
        if ($this->isDuplicate($record)) {
            $problems[] = [
                'field' => null,
                'type' => DirtyDataType::DUPLICATE,
                'level' => DirtyDataLevel::MEDIUM,
                'description' => '数据重复',
                'original_value' => json_encode($record)
            ];
        }
        return $problems;
    }
    /**
     * 获取建议值
     */
    private function getSuggestion($field, $value)
    {
        // 可扩展的自动修复逻辑
        switch ($field) {
            case 'phone':
                if ($value && strlen($value) > 11) {
                    return substr($value, 0, 11);
                }
                break;
            case 'email':
                if ($value && !filter_var($value, FILTER_VALIDATE_EMAIL)) {
                    return preg_replace('/[^a-zA-Z0-9@._-]/', '', $value);
                }
                break;
        }
        return null;
    }
}

批量处理系统

任务队列实现

<?php
class DirtyDataProcessor 
{
    private $db;
    private $redis;
    private $batchSize = 500;
    public function __construct()
    {
        $this->db = new PDO('mysql:host=localhost;dbname=test', 'user', 'pass');
        $this->redis = new Redis();
        $this->redis->connect('127.0.0.1', 6379);
    }
    /**
     * 创建批量处理任务
     */
    public function createBatchTask($table, $conditions = [], $strategy = 'auto_fix')
    {
        $batchNo = date('YmdHis') . '_' . uniqid();
        // 获取符合条件的数据
        $records = $this->getDirtyRecords($table, $conditions);
        // 分批处理
        $chunks = array_chunk($records, $this->batchSize);
        foreach ($chunks as $index => $chunk) {
            $taskData = [
                'batch_no' => $batchNo,
                'table' => $table,
                'records' => $chunk,
                'strategy' => $strategy,
                'chunk_index' => $index,
                'status' => 'pending',
                'created_at' => time()
            ];
            // 放入Redis队列
            $this->redis->rPush('dirty_data:queue', json_encode($taskData));
        }
        return $batchNo;
    }
    /**
     * 处理任务
     */
    public function processTask()
    {
        while ($taskJson = $this->redis->blPop('dirty_data:queue', 5)) {
            $task = json_decode($taskJson[1], true);
            try {
                $task['status'] = 'processing';
                $this->redis->set("dirty_data:task:{$task['batch_no']}:{$task['chunk_index']}", json_encode($task));
                // 执行处理策略
                $this->executeStrategy($task['strategy'], $task['table'], $task['records']);
                $task['status'] = 'completed';
                $task['completed_at'] = time();
                $this->redis->set("dirty_data:task:{$task['batch_no']}:{$task['chunk_index']}", json_encode($task));
            } catch (\Exception $e) {
                $task['status'] = 'failed';
                $task['error'] = $e->getMessage();
                $this->redis->set("dirty_data:task:{$task['batch_no']}:{$task['chunk_index']}", json_encode($task));
                // 记录失败日志
                $this->logFailure($task, $e);
            }
        }
    }
    /**
     * 执行处理策略
     */
    private function executeStrategy($strategy, $table, $records)
    {
        switch ($strategy) {
            case 'auto_fix':
                $this->autoFixStrategy($table, $records);
                break;
            case 'manual_review':
                $this->markForReview($table, $records);
                break;
            case 'delete':
                $this->softDeleteStrategy($table, $records);
                break;
            case 'archive':
                $this->archiveStrategy($table, $records);
                break;
        }
    }
    /**
     * 自动修复策略
     */
    private function autoFixStrategy($table, $records)
    {
        $db = $this->db;
        $db->beginTransaction();
        try {
            $updateStmt = $db->prepare("UPDATE {$table} SET ... WHERE id = ?");
            $markerUpdateStmt = $db->prepare(
                "UPDATE data_quality_markers SET status = 3, handled_at = NOW() WHERE id = ?"
            );
            foreach ($records as $record) {
                // 获取修复建议
                $suggestions = $this->getFixSuggestions($record);
                // 执行修复
                $updateStmt->execute($suggestions['values']);
                // 更新标记状态
                $markerUpdateStmt->execute([$record['marker_id']]);
            }
            $db->commit();
        } catch (\Exception $e) {
            $db->rollBack();
            throw $e;
        }
    }
    /**
     * 获取修复建议
     */
    private function getFixSuggestions($record)
    {
        $stmt = $this->db->prepare(
            "SELECT field_name, suggested_value 
             FROM data_quality_markers 
             WHERE record_id = ? AND status = 0"
        );
        $stmt->execute([$record['id']]);
        $suggestions = $stmt->fetchAll(PDO::FETCH_ASSOC);
        $values = [];
        foreach ($suggestions as $suggestion) {
            if ($suggestion['suggested_value']) {
                $values[$suggestion['field_name']] = $suggestion['suggested_value'];
            }
        }
        return ['values' => $values, 'marker_id' => $record['id']];
    }
}

命令行处理脚本

#!/usr/bin/php
<?php
// batch_process.php
require 'vendor/autoload.php';
class BatchProcessCommand 
{
    private $classifier;
    private $processor;
    public function __construct()
    {
        $this->classifier = new DataQualityClassifier();
        $this->processor = new DirtyDataProcessor();
    }
    /**
     * 全流程处理
     */
    public function run($table, $limit = 1000)
    {
        echo "开始处理数据表: {$table}\n";
        // 1. 获取数据
        echo "正在获取数据...\n";
        $records = $this->fetchRecords($table, $limit);
        if (empty($records)) {
            echo "没有需要处理的数据\n";
            return;
        }
        echo "共获取 {$limit} 条记录\n";
        // 2. 分类标记
        echo "正在进行数据分类标记...\n";
        $batchNo = date('YmdHis');
        $this->classifier->batchMark($table, $records, $batchNo);
        echo "标记完成,批次号: {$batchNo}\n";
        // 3. 统计各类问题
        echo "正在统计问题分布...\n";
        $stats = $this->getStatistics($batchNo);
        $this->printStatistics($stats);
        // 4. 批量处理
        echo "开始批量处理...\n";
        $this->processor->createBatchTask($table, ['batch_no' => $batchNo], 'auto_fix');
        $this->processor->processTask();
        echo "处理完成\n";
        // 5. 生成报告
        $this->generateReport($batchNo);
    }
    /**
     * 获取统计数据
     */
    private function getStatistics($batchNo)
    {
        $stmt = $this->db->prepare(
            "SELECT 
                dirty_type,
                dirty_level,
                COUNT(*) as count
             FROM data_quality_markers
             WHERE batch_no = ?
             GROUP BY dirty_type, dirty_level"
        );
        $stmt->execute([$batchNo]);
        return $stmt->fetchAll(PDO::FETCH_ASSOC);
    }
    /**
     * 打印统计信息
     */
    private function printStatistics($stats)
    {
        echo "\n问题统计:\n";
        echo str_repeat('-', 50) . "\n";
        $typeNames = [
            1 => '数据不完整',
            2 => '数据重复',
            3 => '格式错误',
            4 => '数值超范围',
            5 => '数据不一致',
            6 => '逻辑错误',
            7 => '过期数据',
            8 => '垃圾数据'
        ];
        $levelNames = [
            1 => '致命',
            2 => '高',
            3 => '中',
            4 => '低'
        ];
        foreach ($stats as $stat) {
            $typeName = $typeNames[$stat['dirty_type']] ?? '未知';
            $levelName = $levelNames[$stat['dirty_level']] ?? '未知';
            echo sprintf(
                "%-20s | %-10s | %d条\n",
                $typeName,
                $levelName,
                $stat['count']
            );
        }
        echo str_repeat('-', 50) . "\n";
    }
    /**
     * 生成处理报告
     */
    private function generateReport($batchNo)
    {
        // 生成CSV报告
        $filename = "report_{$batchNo}.csv";
        $fp = fopen($filename, 'w');
        fputcsv($fp, ['ID', '表名', '记录ID', '字段', '问题类型', '级别', '描述', '原始值', '建议值', '状态']);
        $stmt = $this->db->prepare("SELECT * FROM data_quality_markers WHERE batch_no = ?");
        $stmt->execute([$batchNo]);
        while ($row = $stmt->fetch(PDO::FETCH_ASSOC)) {
            fputcsv($fp, $row);
        }
        fclose($fp);
        echo "报告已生成: {$filename}\n";
    }
}
// 执行
$command = new BatchProcessCommand();
$command->run('your_table_name', 1000);

Web管理界面核心功能

<?php
// DirtyDataController.php
class DirtyDataController 
{
    /**
     * 分类展示脏数据
     */
    public function index()
    {
        $type = $_GET['type'] ?? null;
        $level = $_GET['level'] ?? null;
        $status = $_GET['status'] ?? null;
        $query = "SELECT * FROM data_quality_markers WHERE 1=1";
        $params = [];
        if ($type) {
            $query .= " AND dirty_type = ?";
            $params[] = $type;
        }
        if ($level) {
            $query .= " AND dirty_level = ?";
            $params[] = $level;
        }
        if ($status) {
            $query .= " AND status = ?";
            $params[] = $status;
        }
        $query .= " ORDER BY dirty_level ASC, created_at DESC LIMIT 100";
        // 执行查询并渲染视图
    }
    /**
     * 批量操作
     */
    public function batchAction()
    {
        $action = $_POST['action']; // fix, delete, ignore
        $ids = $_POST['ids']; // 选中的记录ID数组
        switch ($action) {
            case 'fix':
                $this->autoFixRecords($ids);
                break;
            case 'delete':
                $this->softDeleteRecords($ids);
                break;
            case 'ignore':
                $this->ignoreRecords($ids);
                break;
        }
        // 返回操作结果
    }
    /**
     * 导出脏数据报告
     */
    public function export()
    {
        $format = $_GET['format'] ?? 'csv';
        $type = $_GET['type'] ?? null;
        // 导出各类格式
        switch ($format) {
            case 'csv':
                $this->exportCSV($type);
                break;
            case 'excel':
                $this->exportExcel($type);
                break;
            case 'json':
                $this->exportJSON($type);
                break;
        }
    }
}

最佳实践建议

数据治理流程

graph LR
    A[数据采集] --> B[实时校验]
    B --> C{校验通过?}
    C -->|是| D[存储]
    C -->|否| E[分类标记]
    E --> F[批量处理队列]
    F --> G[自动修复]
    F --> H[人工审核]
    F --> I[归档删除]
    G --> J[质量检查]
    H --> J
    I --> K[数据回收站]
    J --> L[更新数据状态]

关键优化点

  • 渐进式处理:先处理高严重级别数据
  • 保留审计日志:所有操作可追溯
  • 阈值告警:脏数据占比超过阈值时告警
  • 自动回滚:修复失败时自动回滚
  • 性能优化:使用批量操作减少数据库交互

这套方案已经在多个PHP项目中实践验证,能够有效管理和处理脏数据,提升数据质量,建议根据实际业务场景调整分类规则和处理策略。

抱歉,评论功能暂时关闭!