本文目录导读:

针对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项目中实践验证,能够有效管理和处理脏数据,提升数据质量,建议根据实际业务场景调整分类规则和处理策略。