📖 数据密集型设计

数据生命周期与归档策略

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

一、数据生命周期概述

数据生命周期是指数据从创建到销毁的整个过程。数据生命周期管理能够帮助组织合理利用存储资源、降低成本、确保数据合规性。数据生命周期管理是数据治理的重要组成部分。

二、数据生命周期阶段

2.1 数据生命周期阶段表

阶段 描述 存储需求 访问频率 管理策略
创建 数据生成和采集 高性能存储 实时写入
活跃 数据频繁访问和修改 高性能存储 实时读写
非活跃 数据访问频率降低 中等性能存储 定期访问
归档 数据长期保存但很少访问 低成本存储 偶尔访问
销毁 数据不再需要 安全删除

2.2 数据生命周期流程图

flowchart TD A[数据创建] --> B[数据活跃] B --> C[数据非活跃] C --> D[数据归档] D --> E[数据销毁] B -->|频繁访问| B C -->|访问频率降低| C D -->|合规要求| D F[存储层] --> F1[SSD/高性能存储] F --> F2[SATA/中等存储] F --> F3[对象存储/归档存储] A & B --> F1 C --> F2 D --> F3

三、冷热分离策略

3.1 冷热分离原理

冷热分离是根据数据访问频率将数据存储在不同类型的存储介质上:

graph TD A[数据访问层] --> B{数据类型} B -->|热数据| C[热存储 - SSD] B -->|温数据| D[温存储 - SATA] B -->|冷数据| E[冷存储 - 对象存储] C --> C1[高IOPS] C --> C2[低延迟] C --> C3[高成本] D --> D1[中等IOPS] D --> D2[中等延迟] D --> D3[中等成本] E --> E1[低IOPS] E --> E2[高延迟] E --> E3[低成本]

3.2 冷热数据识别

public class HotColdDataClassifier
{
    public async Task<DataClassificationResult> ClassifyAsync(string tableName)
    {
        var accessStats = await _accessMonitor.GetAccessStatisticsAsync(tableName);
        
        var hotThreshold = TimeSpan.FromDays(30);
        var warmThreshold = TimeSpan.FromDays(90);
        
        var hotData = accessStats.Where(s => s.LastAccessTime >= DateTime.Now - hotThreshold).ToList();
        var warmData = accessStats.Where(s => 
            s.LastAccessTime >= DateTime.Now - warmThreshold && 
            s.LastAccessTime < DateTime.Now - hotThreshold).ToList();
        var coldData = accessStats.Where(s => s.LastAccessTime < DateTime.Now - warmThreshold).ToList();
        
        return new DataClassificationResult
        {
            TableName = tableName,
            HotDataPercentage = hotData.Count / (double)accessStats.Count * 100,
            WarmDataPercentage = warmData.Count / (double)accessStats.Count * 100,
            ColdDataPercentage = coldData.Count / (double)accessStats.Count * 100,
            HotDataSize = hotData.Sum(s => s.Size),
            WarmDataSize = warmData.Sum(s => s.Size),
            ColdDataSize = coldData.Sum(s => s.Size)
        };
    }
    
    public async Task<List<DataPartition>> GetPartitionRecommendationsAsync(string tableName)
    {
        var classification = await ClassifyAsync(tableName);
        
        var recommendations = new List<DataPartition>();
        
        if (classification.ColdDataPercentage > 50)
        {
            recommendations.Add(new DataPartition
            {
                PartitionName = "cold_data",
                StorageType = StorageType.Archive,
                DataRange = $"CreatedAt < '{DateTime.Now.AddDays(-90):yyyy-MM-dd}'"
            });
        }
        
        if (classification.WarmDataPercentage > 30)
        {
            recommendations.Add(new DataPartition
            {
                PartitionName = "warm_data",
                StorageType = StorageType.Warm,
                DataRange = $"CreatedAt >= '{DateTime.Now.AddDays(-90):yyyy-MM-dd}' AND CreatedAt < '{DateTime.Now.AddDays(-30):yyyy-MM-dd}'"
            });
        }
        
        return recommendations;
    }
}

3.3 MySQL冷热分离实现

public class MysqlHotColdSeparationService
{
    public async Task ImplementHotColdSeparationAsync(string tableName)
    {
        await CreateColdTableAsync(tableName);
        await MigrateColdDataAsync(tableName);
        await CreatePartitionTriggerAsync(tableName);
    }
    
    private async Task CreateColdTableAsync(string tableName)
    {
        var coldTableName = $"{tableName}_cold";
        
        var sql = @$"
            CREATE TABLE {coldTableName} LIKE {tableName};
            ALTER TABLE {coldTableName} ENGINE = ARCHIVE;
        ";
        
        await _dbContext.Database.ExecuteSqlRawAsync(sql);
    }
    
    private async Task MigrateColdDataAsync(string tableName)
    {
        var coldTableName = $"{tableName}_cold";
        
        var sql = @$"
            INSERT INTO {coldTableName} 
            SELECT * FROM {tableName} 
            WHERE created_at < DATE_SUB(NOW(), INTERVAL 90 DAY);
            
            DELETE FROM {tableName} 
            WHERE created_at < DATE_SUB(NOW(), INTERVAL 90 DAY);
        ";
        
        await _dbContext.Database.ExecuteSqlRawAsync(sql);
    }
    
    private async Task CreatePartitionTriggerAsync(string tableName)
    {
        var coldTableName = $"{tableName}_cold";
        
        var sql = @$"
            CREATE TRIGGER migrate_to_cold 
            BEFORE INSERT ON {tableName}
            FOR EACH ROW
            BEGIN
                IF NEW.created_at < DATE_SUB(NOW(), INTERVAL 90 DAY) THEN
                    INSERT INTO {coldTableName} VALUES (NEW.*);
                    SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = 'Data should go to cold table';
                END IF;
            END;
        ";
        
        await _dbContext.Database.ExecuteSqlRawAsync(sql);
    }
    
    public async Task QueryHotAndColdAsync(string tableName, string query)
    {
        var coldTableName = $"{tableName}_cold";
        
        var sql = @$"
            SELECT * FROM {tableName} WHERE {query}
            UNION ALL
            SELECT * FROM {coldTableName} WHERE {query};
        ";
        
        await _dbContext.QueryAsync(sql);
    }
}

四、数据归档策略

四、数据归档策略

4.1 归档策略对比

策略 描述 优点 缺点 适用场景
数据库归档 数据仍在数据库中,存储在归档表 查询方便 成本较高 需要偶尔查询
文件归档 数据导出到文件 成本低 查询不便 很少查询
对象存储归档 数据存储到对象存储 成本极低 延迟高 长期归档
数据湖归档 数据存储到数据湖 支持分析 复杂度高 需要数据分析

4.2 对象存储归档

public class ObjectStorageArchivalService
{
    private readonly IBlobStorage _blobStorage;
    
    public ObjectStorageArchivalService(IBlobStorage blobStorage)
    {
        _blobStorage = blobStorage;
    }
    
    public async Task ArchiveToObjectStorageAsync(string tableName, DateTime cutoffDate)
    {
        var data = await FetchDataForArchivalAsync(tableName, cutoffDate);
        
        var parquetService = new ParquetStorageService();
        var tempFilePath = Path.GetTempFileName();
        
        await parquetService.WriteParquetAsync(tempFilePath, data);
        
        var archivePath = GenerateArchivePath(tableName, cutoffDate);
        
        await _blobStorage.UploadFileAsync(tempFilePath, archivePath);
        
        await DeleteArchivedDataAsync(tableName, cutoffDate);
        
        await LogArchivalAsync(tableName, cutoffDate, data.Count);
        
        File.Delete(tempFilePath);
    }
    
    private async Task<List<dynamic>> FetchDataForArchivalAsync(string tableName, DateTime cutoffDate)
    {
        var sql = $"SELECT * FROM {tableName} WHERE created_at < @CutoffDate";
        return await _dbContext.QueryAsync(sql, new { CutoffDate = cutoffDate });
    }
    
    private string GenerateArchivePath(string tableName, DateTime cutoffDate)
    {
        return $"archive/{tableName}/{cutoffDate:yyyy}/{cutoffDate:MM}/{cutoffDate:dd}/data.parquet";
    }
    
    private async Task DeleteArchivedDataAsync(string tableName, DateTime cutoffDate)
    {
        var sql = $"DELETE FROM {tableName} WHERE created_at < @CutoffDate";
        await _dbContext.Database.ExecuteSqlRawAsync(sql, new { CutoffDate = cutoffDate });
    }
    
    private async Task LogArchivalAsync(string tableName, DateTime cutoffDate, int recordCount)
    {
        await _archivalRepository.CreateLogAsync(new ArchivalLog
        {
            TableName = tableName,
            CutoffDate = cutoffDate,
            RecordCount = recordCount,
            ArchiveTime = DateTime.Now,
            ArchiveType = ArchiveType.ObjectStorage
        });
    }
    
    public async Task<List<dynamic>> RetrieveFromArchiveAsync(string tableName, DateTime date)
    {
        var archivePath = GenerateArchivePath(tableName, date);
        
        var tempFilePath = Path.GetTempFileName();
        
        await _blobStorage.DownloadFileAsync(archivePath, tempFilePath);
        
        var parquetService = new ParquetStorageService();
        var data = await parquetService.ReadParquetAsync(tempFilePath);
        
        File.Delete(tempFilePath);
        
        return data;
    }
}

4.3 数据湖归档

public class DataLakeArchivalService
{
    public async Task ArchiveToDataLakeAsync(string tableName, DateTime cutoffDate)
    {
        var data = await FetchDataForArchivalAsync(tableName, cutoffDate);
        
        var deltaTablePath = GetDeltaTablePath(tableName);
        
        await WriteToDeltaLakeAsync(deltaTablePath, data);
        
        await DeleteArchivedDataAsync(tableName, cutoffDate);
        
        await UpdatePartitionStatisticsAsync(deltaTablePath);
    }
    
    private string GetDeltaTablePath(string tableName)
    {
        return $"s3://datalake/archive/{tableName}";
    }
    
    private async Task WriteToDeltaLakeAsync(string path, List<dynamic> data)
    {
        var spark = SparkSession.Builder()
            .AppName("DataLakeArchival")
            .GetOrCreate();
        
        var df = spark.CreateDataFrame(data);
        
        df.Write()
            .Format("delta")
            .Mode("append")
            .PartitionBy("year", "month", "day")
            .Save(path);
    }
    
    private async Task UpdatePartitionStatisticsAsync(string path)
    {
        var spark = SparkSession.Builder()
            .AppName("UpdateStatistics")
            .GetOrCreate();
        
        spark.Sql($"ANALYZE TABLE delta.`{path}` COMPUTE STATISTICS");
    }
    
    public async Task<List<dynamic>> QueryFromDataLakeAsync(string tableName, string query)
    {
        var deltaTablePath = GetDeltaTablePath(tableName);
        
        var spark = SparkSession.Builder()
            .AppName("DataLakeQuery")
            .GetOrCreate();
        
        var df = spark.Read().Format("delta").Load(deltaTablePath);
        
        var results = df.Sql(query).CollectAsList();
        
        return results.Select(r => r.AsDynamic()).ToList();
    }
}

五、数据过期策略

5.1 过期策略类型

策略 描述 适用场景
时间过期 按时间自动删除 日志、会话
大小过期 超过大小限制后删除 缓存、临时文件
版本过期 保留指定版本数 版本控制
合规过期 按法规要求删除 用户数据、合规数据

5.2 Redis过期策略

public class RedisExpirationService
{
    public async Task SetWithExpirationAsync(string key, object value, TimeSpan expiration)
    {
        await _cache.SetStringAsync(key, JsonSerializer.Serialize(value), 
            new DistributedCacheEntryOptions
            {
                AbsoluteExpirationRelativeToNow = expiration
            });
    }
    
    public async Task SetSlidingExpirationAsync(string key, object value, TimeSpan slidingExpiration)
    {
        await _cache.SetStringAsync(key, JsonSerializer.Serialize(value), 
            new DistributedCacheEntryOptions
            {
                SlidingExpiration = slidingExpiration
            });
    }
    
    public async Task SetAbsoluteExpirationAsync(string key, object value, DateTime absoluteExpiration)
    {
        await _cache.SetStringAsync(key, JsonSerializer.Serialize(value), 
            new DistributedCacheEntryOptions
            {
                AbsoluteExpiration = absoluteExpiration
            });
    }
    
    public async Task SetWithComplexExpirationAsync(string key, object value, 
        TimeSpan absoluteExpiration, TimeSpan slidingExpiration)
    {
        await _cache.SetStringAsync(key, JsonSerializer.Serialize(value), 
            new DistributedCacheEntryOptions
            {
                AbsoluteExpirationRelativeToNow = absoluteExpiration,
                SlidingExpiration = slidingExpiration
            });
    }
    
    public async Task<long> CleanupExpiredKeysAsync(string pattern)
    {
        var keys = await _cache.SearchKeysAsync(pattern);
        var deletedCount = 0;
        
        foreach (var key in keys)
        {
            var ttl = await _cache.GetTtlAsync(key);
            
            if (ttl.HasValue && ttl.Value <= TimeSpan.Zero)
            {
                await _cache.RemoveAsync(key);
                deletedCount++;
            }
        }
        
        return deletedCount;
    }
}

5.3 MongoDB过期策略

public class MongoExpirationService
{
    public async Task CreateTTLIndexAsync(string collectionName, string fieldName, int expireAfterSeconds)
    {
        var collection = _mongoDatabase.GetCollection<BsonDocument>(collectionName);
        
        var indexKeys = Builders<BsonDocument>.IndexKeys.Ascending(fieldName);
        var indexOptions = new CreateIndexOptions { ExpireAfter = TimeSpan.FromSeconds(expireAfterSeconds) };
        
        await collection.Indexes.CreateOneAsync(new CreateIndexModel<BsonDocument>(indexKeys, indexOptions));
    }
    
    public async Task SetDocumentExpirationAsync<T>(string collectionName, T document, DateTime expirationTime)
    {
        var collection = _mongoDatabase.GetCollection<T>(collectionName);
        
        var bsonDocument = document.ToBsonDocument();
        bsonDocument["expireAt"] = expirationTime;
        
        await collection.InsertOneAsync(BsonSerializer.Deserialize<T>(bsonDocument));
    }
    
    public async Task<long> DeleteExpiredDocumentsAsync(string collectionName)
    {
        var collection = _mongoDatabase.GetCollection<BsonDocument>(collectionName);
        
        var filter = Builders<BsonDocument>.Filter.Lt("expireAt", DateTime.Now);
        var result = await collection.DeleteManyAsync(filter);
        
        return result.DeletedCount;
    }
}

六、数据销毁

6.1 数据销毁方法

方法 描述 安全性 适用场景
逻辑删除 标记删除,数据仍存在 可恢复删除
物理删除 直接删除数据 一般数据
加密删除 加密数据后删除密钥 敏感数据
安全擦除 多次覆写数据 极高 存储设备销毁

6.2 安全数据销毁

public class SecureDataDestructionService
{
    public async Task LogicallyDeleteAsync(string tableName, int id)
    {
        var sql = $"UPDATE {tableName} SET is_deleted = 1, deleted_at = NOW() WHERE id = @Id";
        await _dbContext.Database.ExecuteSqlRawAsync(sql, new { Id = id });
    }
    
    public async Task PhysicallyDeleteAsync(string tableName, int id)
    {
        var sql = $"DELETE FROM {tableName} WHERE id = @Id";
        await _dbContext.Database.ExecuteSqlRawAsync(sql, new { Id = id });
    }
    
    public async Task SecureDeleteAsync(string tableName, int id)
    {
        var data = await FetchDataAsync(tableName, id);
        
        var encryptedData = _encryptionService.Encrypt(data);
        
        var sql = $"UPDATE {tableName} SET data = @EncryptedData, is_encrypted = 1 WHERE id = @Id";
        await _dbContext.Database.ExecuteSqlRawAsync(sql, new { EncryptedData = encryptedData, Id = id });
        
        await _keyRepository.RevokeKeyAsync(GetDataKeyId(id));
        
        await _auditLogger.LogDataDestruction(tableName, id, DestructionType.SecureDelete);
    }
    
    public async Task PurgeDeletedDataAsync(string tableName)
    {
        var sql = $"DELETE FROM {tableName} WHERE is_deleted = 1 AND deleted_at < DATE_SUB(NOW(), INTERVAL 30 DAY)";
        await _dbContext.Database.ExecuteSqlRawAsync(sql);
    }
    
    public async Task DestroyUserDataAsync(int userId)
    {
        var userTables = await _metadataRepository.GetUserRelatedTablesAsync(userId);
        
        foreach (var table in userTables)
        {
            await SecureDeleteAsync(table.TableName, userId);
        }
        
        await _auditLogger.LogUserDataDestruction(userId);
    }
}

七、数据生命周期管理流程

7.1 生命周期管理流程

flowchart TD A[数据创建] --> B[数据使用] B --> C[访问频率分析] C --> D{访问频率} D -->|高| B D -->|中| E[移至温存储] D -->|低| F[移至冷存储] E --> G[定期访问检查] G -->|仍中| E G -->|变高| B G -->|变低| F F --> H[合规检查] H -->|需要保留| I[长期归档] H -->|无需保留| J[数据销毁] I --> K[归档管理] J --> L[销毁审计]

7.2 生命周期管理自动化

public class DataLifecycleManager
{
    public async Task ManageLifecycleAsync()
    {
        await ClassifyDataAsync();
        await MigrateHotToWarmAsync();
        await MigrateWarmToColdAsync();
        await ArchiveColdDataAsync();
        await DeleteExpiredDataAsync();
    }
    
    private async Task ClassifyDataAsync()
    {
        var tables = await _metadataRepository.GetTablesAsync();
        
        foreach (var table in tables)
        {
            var classification = await _hotColdClassifier.ClassifyAsync(table.Name);
            await _lifecycleRepository.UpdateClassificationAsync(table.Name, classification);
        }
    }
    
    private async Task MigrateHotToWarmAsync()
    {
        var tables = await _lifecycleRepository.GetTablesForMigrationAsync(
            LifecyclePhase.Hot, LifecyclePhase.Warm);
        
        foreach (var table in tables)
        {
            await _migrationService.MigrateAsync(table.Name, StorageType.Warm);
        }
    }
    
    private async Task MigrateWarmToColdAsync()
    {
        var tables = await _lifecycleRepository.GetTablesForMigrationAsync(
            LifecyclePhase.Warm, LifecyclePhase.Cold);
        
        foreach (var table in tables)
        {
            await _migrationService.MigrateAsync(table.Name, StorageType.Archive);
        }
    }
    
    private async Task ArchiveColdDataAsync()
    {
        var tables = await _lifecycleRepository.GetTablesForArchivalAsync();
        
        foreach (var table in tables)
        {
            await _archivalService.ArchiveToObjectStorageAsync(table.Name, GetCutoffDate(table));
        }
    }
    
    private async Task DeleteExpiredDataAsync()
    {
        var tables = await _lifecycleRepository.GetTablesWithExpiredDataAsync();
        
        foreach (var table in tables)
        {
            await _expirationService.DeleteExpiredDocumentsAsync(table.Name);
        }
    }
    
    private DateTime GetCutoffDate(TableMetadata table)
    {
        return table.RetentionPolicy switch
        {
            RetentionPolicy.ShortTerm => DateTime.Now.AddMonths(-3),
            RetentionPolicy.MediumTerm => DateTime.Now.AddMonths(-12),
            RetentionPolicy.LongTerm => DateTime.Now.AddYears(-5),
            _ => DateTime.Now.AddMonths(-3)
        };
    }
}

八、数据生命周期监控

8.1 监控指标

public class DataLifecycleMetrics
{
    public string TableName { get; set; }
    public LifecyclePhase CurrentPhase { get; set; }
    public long HotDataSize { get; set; }
    public long WarmDataSize { get; set; }
    public long ColdDataSize { get; set; }
    public double StorageCost { get; set; }
    public int ArchiveCount { get; set; }
    public int DeleteCount { get; set; }
}

public class DataLifecycleMonitor
{
    public async Task<DataLifecycleMetrics> GetMetricsAsync(string tableName)
    {
        var classification = await _hotColdClassifier.ClassifyAsync(tableName);
        
        return new DataLifecycleMetrics
        {
            TableName = tableName,
            CurrentPhase = classification.ColdDataPercentage > 50 
                ? LifecyclePhase.Cold 
                : classification.WarmDataPercentage > 30 
                    ? LifecyclePhase.Warm 
                    : LifecyclePhase.Hot,
            HotDataSize = classification.HotDataSize,
            WarmDataSize = classification.WarmDataSize,
            ColdDataSize = classification.ColdDataSize,
            StorageCost = await CalculateStorageCostAsync(classification),
            ArchiveCount = await _archivalRepository.GetArchiveCountAsync(tableName),
            DeleteCount = await _expirationRepository.GetDeleteCountAsync(tableName)
        };
    }
    
    private async Task<double> CalculateStorageCostAsync(DataClassificationResult classification)
    {
        var hotCost = classification.HotDataSize * 0.02;
        var warmCost = classification.WarmDataSize * 0.01;
        var coldCost = classification.ColdDataSize * 0.001;
        
        return hotCost + warmCost + coldCost;
    }
}

九、数据生命周期最佳实践

9.1 制定生命周期策略

根据业务需求和合规要求,制定数据生命周期策略。

9.2 自动化管理

自动化数据生命周期管理,减少人工干预。

9.3 成本优化

通过冷热分离和归档策略,降低存储成本。

9.4 合规性保障

确保数据销毁符合法规要求,保留审计日志。

9.5 可恢复性

确保归档数据可以快速恢复。

十、总结

数据生命周期管理是数据密集型应用中重要的成本优化和合规保障手段。通过冷热分离策略、数据归档和过期管理,能够显著降低存储成本。数据销毁需要符合法规要求,确保数据安全。自动化的数据生命周期管理能够提高效率,减少人工干预。合理的数据生命周期管理策略能够在满足业务需求的同时,最大化数据价值和最小化成本。