📖 数据密集型设计

数据归档与生命周期管理

深入探讨数据归档策略与生命周期管理

一、数据生命周期概述

数据生命周期管理是对数据从创建到销毁的全过程管理,在数据密集型应用中,有效的数据生命周期管理能够显著降低存储成本、提高系统性能。

二、数据生命周期阶段

2.1 数据生命周期阶段

graph TD A[数据创建] --> B[活跃期] B --> C[过渡期] C --> D[归档期] D --> E[销毁期] B --> B1[高频访问] B1 --> B2[高性能存储] C --> C1[中频访问] C1 --> C2[中等性能存储] D --> D1[低频访问] D1 --> D2[低成本存储] E --> E1[合规要求] E1 --> E2[安全销毁]

2.2 数据生命周期阶段对比

阶段 描述 访问频率 存储类型 保留时间
创建期 数据刚创建 内存/SSD 即时
活跃期 频繁访问 SSD/HDD 0-90天
过渡期 中等访问 HDD 90-365天
归档期 低频访问 对象存储 1-7年
销毁期 不再需要 安全销毁 合规要求

三、数据归档策略

3.1 归档策略对比

策略 描述 优点 缺点 适用场景
时间驱动 按时间归档 简单 不考虑访问频率 日志数据
访问频率 按访问频率归档 精准 需要追踪 文档数据
大小驱动 按大小归档 控制空间 可能提前归档 文件存储
综合策略 多种条件组合 灵活 复杂 企业级应用

3.2 归档策略实现

public class DataArchivalService
{
    public async Task ArchiveDataAsync(DataArchivalRequest request)
    {
        var data = await _database.GetDataAsync(request.TableName, request.Filter);
        
        var archivalDestination = SelectArchivalDestination(request.Policy);
        
        await _archivalStorage.StoreAsync(archivalDestination, data);
        
        await _database.MarkAsArchivedAsync(request.TableName, request.Filter);
    }
    
    private string SelectArchivalDestination(ArchivalPolicy policy)
    {
        return policy.DestinationType switch
        {
            ArchivalDestination.ObjectStorage => "s3://archive",
            ArchivalDestination.Hadoop => "hdfs://archive",
            ArchivalDestination.ColdStorage => "cold://archive",
            _ => "s3://archive"
        };
    }
    
    public async Task RestoreDataAsync(DataRestoreRequest request)
    {
        var data = await _archivalStorage.GetAsync(request.Location);
        
        await _database.RestoreDataAsync(request.TableName, data);
        
        await _archivalStorage.DeleteAsync(request.Location);
    }
    
    public async Task PurgeExpiredDataAsync(string tableName, TimeSpan retentionPeriod)
    {
        var expiredData = await _database.GetExpiredDataAsync(tableName, retentionPeriod);
        
        await _database.DeleteDataAsync(tableName, expiredData);
        
        await _auditService.LogPurgeAsync(tableName, expiredData.Count);
    }
}

public class DataArchivalRequest
{
    public string TableName { get; set; }
    public string Filter { get; set; }
    public ArchivalPolicy Policy { get; set; }
}

public class ArchivalPolicy
{
    public string Name { get; set; }
    public ArchivalDestination DestinationType { get; set; }
    public TimeSpan RetentionPeriod { get; set; }
    public CompressionAlgorithm CompressionAlgorithm { get; set; }
}

public enum ArchivalDestination { ObjectStorage, Hadoop, ColdStorage, Tape }

四、冷热分离

4.1 冷热分离架构

graph TD A[数据访问] --> B[路由层] B --> C{数据类型?} C -->|热数据| D[热存储层] C -->|冷数据| E[冷存储层] D --> D1[SSD存储] D1 --> D2[高频访问] E --> E1[对象存储] E1 --> E2[低频访问] F[数据迁移] --> G[热→冷] G --> H[定时任务] F --> I[冷→热] I --> J[按需恢复]

4.2 冷热分离实现

public class HotColdSeparationService
{
    public async Task GetDataAsync(string key)
    {
        var dataLocation = await _metadataService.GetDataLocationAsync(key);
        
        switch (dataLocation.LocationType)
        {
            case DataLocationType.Hot:
                return await _hotStorage.GetAsync(key);
            case DataLocationType.Cold:
                return await _coldStorage.GetAsync(key);
            case DataLocationType.Archived:
                return await RestoreAndGetAsync(key);
            default:
                throw new KeyNotFoundException("Data not found");
        }
    }
    
    public async Task StoreDataAsync(string key, object data, DataLocationType locationType)
    {
        switch (locationType)
        {
            case DataLocationType.Hot:
                await _hotStorage.PutAsync(key, data);
                break;
            case DataLocationType.Cold:
                await _coldStorage.PutAsync(key, data);
                break;
        }
        
        await _metadataService.SetDataLocationAsync(key, locationType);
    }
    
    public async Task MigrateToColdAsync(string key)
    {
        var data = await _hotStorage.GetAsync(key);
        
        await _coldStorage.PutAsync(key, data);
        await _hotStorage.DeleteAsync(key);
        
        await _metadataService.SetDataLocationAsync(key, DataLocationType.Cold);
    }
    
    public async Task MigrateToHotAsync(string key)
    {
        var data = await _coldStorage.GetAsync(key);
        
        await _hotStorage.PutAsync(key, data);
        await _coldStorage.DeleteAsync(key);
        
        await _metadataService.SetDataLocationAsync(key, DataLocationType.Hot);
    }
    
    private async Task RestoreAndGetAsync(string key)
    {
        var data = await _archivalStorage.GetAsync(key);
        
        await _hotStorage.PutAsync(key, data);
        await _metadataService.SetDataLocationAsync(key, DataLocationType.Hot);
        
        return data;
    }
}

public enum DataLocationType { Hot, Cold, Archived }
            
            

4.3 冷热数据识别

public class HotColdDataIdentifier
{
    public async Task> IdentifyHotDataAsync(string tableName, TimeSpan timeWindow)
    {
        var accessLogs = await _accessLogService.GetAccessLogsAsync(tableName, timeWindow);
        
        var hotKeys = accessLogs
            .GroupBy(log => log.Key)
            .Where(g => g.Count() > 100)
            .Select(g => g.Key)
            .ToList();
        
        return hotKeys;
    }
    
    public async Task> IdentifyColdDataAsync(string tableName, TimeSpan timeWindow)
    {
        var accessLogs = await _accessLogService.GetAccessLogsAsync(tableName, timeWindow);
        
        var coldKeys = accessLogs
            .GroupBy(log => log.Key)
            .Where(g => g.Count() <= 10)
            .Select(g => g.Key)
            .ToList();
        
        return coldKeys;
    }
    
    public async Task> IdentifyStaleDataAsync(string tableName, TimeSpan stalePeriod)
    {
        var dataMetadata = await _metadataService.GetDataMetadataAsync(tableName);
        
        var staleKeys = dataMetadata
            .Where(m => m.LastAccessed < DateTime.UtcNow - stalePeriod)
            .Select(m => m.Key)
            .ToList();
        
        return staleKeys;
    }
    
    public async Task ExecuteHotColdMigrationAsync(string tableName)
    {
        var coldData = await IdentifyColdDataAsync(tableName, TimeSpan.FromDays(30));
        var staleData = await IdentifyStaleDataAsync(tableName, TimeSpan.FromDays(90));
        
        var dataToMigrate = coldData.Union(staleData).ToList();
        
        foreach (var key in dataToMigrate)
        {
            await _hotColdService.MigrateToColdAsync(key);
        }
        
        await _auditService.LogMigrationAsync(tableName, dataToMigrate.Count);
    }
}

五、数据过期与清理

5.1 数据过期策略

public class DataExpirationService
{
    public async Task CleanupExpiredDataAsync(DataExpirationOptions options)
    {
        foreach (var table in options.Tables)
        {
            await CleanupTableAsync(table, options.RetentionPeriod);
        }
    }
    
    private async Task CleanupTableAsync(string tableName, TimeSpan retentionPeriod)
    {
        var expiredData = await _database.GetExpiredDataAsync(tableName, retentionPeriod);
        
        foreach (var data in expiredData)
        {
            await ProcessExpiredDataAsync(tableName, data);
        }
    }
    
    private async Task ProcessExpiredDataAsync(string tableName, ExpiredData data)
    {
        switch (data.ExpirationAction)
        {
            case ExpirationAction.Delete:
                await _database.DeleteDataAsync(tableName, data.Id);
                break;
            case ExpirationAction.Archive:
                await _archivalService.ArchiveDataAsync(new DataArchivalRequest
                {
                    TableName = tableName,
                    Filter = $"Id = {data.Id}",
                    Policy = data.ArchivalPolicy
                });
                await _database.DeleteDataAsync(tableName, data.Id);
                break;
            case ExpirationAction.Anonymize:
                await _anonymizationService.AnonymizeAsync(tableName, data.Id);
                break;
        }
        
        await _auditService.LogExpirationAsync(tableName, data.Id, data.ExpirationAction);
    }
}

public class DataExpirationOptions
{
    public List Tables { get; set; }
    public TimeSpan RetentionPeriod { get; set; }
    public ExpirationAction DefaultAction { get; set; }
}

public enum ExpirationAction { Delete, Archive, Anonymize }

5.2 TTL索引实现

public class TtlIndexService
{
    public void CreateTtlIndex(string tableName, string columnName, TimeSpan ttl)
    {
        _database.ExecuteSql($@"
            CREATE INDEX IX_{tableName}_{columnName}_TTL 
            ON {tableName}({columnName})
            WHERE {columnName} IS NOT NULL
        ");
        
        _ttlMetadataService.RegisterTtlPolicy(new TtlPolicy
        {
            TableName = tableName,
            ColumnName = columnName,
            Ttl = ttl,
            Enabled = true
        });
    }
    
    public async Task ProcessTtlExpirationsAsync()
    {
        var policies = await _ttlMetadataService.GetEnabledPoliciesAsync();
        
        foreach (var policy in policies)
        {
            await ProcessTtlPolicyAsync(policy);
        }
    }
    
    private async Task ProcessTtlPolicyAsync(TtlPolicy policy)
    {
        var expirationTime = DateTime.UtcNow - policy.Ttl;
        
        var expiredData = await _database.GetDataAsync($@"
            SELECT Id FROM {policy.TableName} 
            WHERE {policy.ColumnName} < '{expirationTime:yyyy-MM-dd HH:mm:ss}'
        ");
        
        foreach (var data in expiredData)
        {
            await _database.DeleteDataAsync(policy.TableName, data.Id);
        }
        
        await _auditService.LogTtlCleanupAsync(policy.TableName, expiredData.Count);
    }
}

public class TtlPolicy
{
    public string TableName { get; set; }
    public string ColumnName { get; set; }
    public TimeSpan Ttl { get; set; }
    public bool Enabled { get; set; }
}

六、数据归档监控与审计

6.1 归档监控

public class DataArchivalMonitor
{
    public async Task GetMetricsAsync()
    {
        var metrics = new ArchivalMetrics();
        
        var policies = await _archivalPolicyService.GetAllPoliciesAsync();
        
        foreach (var policy in policies)
        {
            var policyMetrics = await _archivalStorage.GetPolicyMetricsAsync(policy.Name);
            
            metrics.TotalArchivedBytes += policyMetrics.ArchivedBytes;
            metrics.TotalArchivedRecords += policyMetrics.ArchivedRecords;
            metrics.TotalPolicies++;
        }
        
        metrics.ArchivedTables = await _database.GetArchivedTableCountAsync();
        metrics.ActiveTables = await _database.GetActiveTableCountAsync();
        
        return metrics;
    }
    
    public async Task MonitorAsync()
    {
        var metrics = await GetMetricsAsync();
        
        if (metrics.TotalArchivedBytes > 100 * 1024 * 1024 * 1024)
        {
            await _alertService.SendAlert("归档数据过大", 
                $"归档数据大小: {metrics.TotalArchivedBytes / 1024 / 1024 / 1024} GB");
        }
    }
}

public class ArchivalMetrics
{
    public long TotalArchivedBytes { get; set; }
    public long TotalArchivedRecords { get; set; }
    public int TotalPolicies { get; set; }
    public int ArchivedTables { get; set; }
    public int ActiveTables { get; set; }
}

6.2 归档审计

public class DataArchivalAuditService
{
    public async Task LogArchivalAsync(string tableName, int recordCount, string destination)
    {
        var auditRecord = new ArchivalAuditRecord
        {
            Id = Guid.NewGuid().ToString(),
            TableName = tableName,
            RecordCount = recordCount,
            Destination = destination,
            Timestamp = DateTime.UtcNow,
            Action = ArchivalAction.Archive
        };
        
        await _auditRepository.CreateAsync(auditRecord);
    }
    
    public async Task LogRestoreAsync(string tableName, int recordCount, string source)
    {
        var auditRecord = new ArchivalAuditRecord
        {
            Id = Guid.NewGuid().ToString(),
            TableName = tableName,
            RecordCount = recordCount,
            Destination = source,
            Timestamp = DateTime.UtcNow,
            Action = ArchivalAction.Restore
        };
        
        await _auditRepository.CreateAsync(auditRecord);
    }
    
    public async Task LogExpirationAsync(string tableName, int recordCount, ExpirationAction action)
    {
        var auditRecord = new ArchivalAuditRecord
        {
            Id = Guid.NewGuid().ToString(),
            TableName = tableName,
            RecordCount = recordCount,
            Timestamp = DateTime.UtcNow,
            Action = ArchivalAction.Expiration,
            ExpirationAction = action
        };
        
        await _auditRepository.CreateAsync(auditRecord);
    }
    
    public async Task> GetAuditLogsAsync(string tableName, DateTime? startTime, DateTime? endTime)
    {
        return await _auditRepository.GetLogsAsync(tableName, startTime, endTime);
    }
}

public class ArchivalAuditRecord
{
    public string Id { get; set; }
    public string TableName { get; set; }
    public int RecordCount { get; set; }
    public string Destination { get; set; }
    public DateTime Timestamp { get; set; }
    public ArchivalAction Action { get; set; }
    public ExpirationAction? ExpirationAction { get; set; }
}

public enum ArchivalAction { Archive, Restore, Expiration }

七、数据生命周期管理最佳实践

7.1 生命周期策略设计

数据类型 活跃期 过渡期 归档期 销毁期
交易数据 90天 365天 7年 7年后
日志数据 30天 90天 1年 1年后
用户数据 无限 - - 用户注销后
临时数据 7天 - - 7天后

7.2 归档实施建议

  • 制定数据生命周期策略
  • 实施冷热分离
  • 配置TTL索引
  • 做好归档监控
  • 保留归档审计日志

7.3 合规检查清单

public class DataLifecycleComplianceChecker
{
    public async Task CheckComplianceAsync()
    {
        var report = new ComplianceReport();
        
        report.TtlPoliciesCount = await _ttlService.GetPolicyCountAsync();
        report.ArchivalPoliciesCount = await _archivalService.GetPolicyCountAsync();
        report.AuditLogsCount = await _auditService.GetLogsCountAsync();
        
        var expiredData = await _database.GetExpiredDataCountAsync();
        report.OverdueDataCount = expiredData;
        
        report.IsCompliant = report.TtlPoliciesCount > 0 && 
                           report.ArchivalPoliciesCount > 0 && 
                           report.AuditLogsCount > 0 &&
                           report.OverdueDataCount == 0;
        
        return report;
    }
}

八、总结

数据归档与生命周期管理是降低存储成本、提高系统性能的关键。通过合理制定生命周期策略、实施冷热分离、配置TTL索引,能够构建高效的数据生命周期管理系统。做好监控和审计,保障数据合规性,能够避免法律风险。