📖 数据密集型设计

数据压缩算法与存储优化

深入探讨压缩算法原理与存储优化策略

一、数据压缩概述

数据压缩是将数据体积减小的技术,在数据密集型应用中,数据压缩能够显著减少存储空间和网络传输开销,提高系统性能。

二、压缩算法分类

2.1 压缩算法分类

graph TD A[压缩算法] --> B[有损压缩] A --> C[无损压缩] B --> B1[图像压缩] B --> B2[音频压缩] B --> B3[视频压缩] C --> C1[通用压缩] C1 --> C11[LZ77] C1 --> C12[LZ78] C1 --> C13[Huffman] C1 --> C14[LZW] C --> C2[专用压缩] C2 --> C21[列式压缩] C2 --> C22[字典编码] C2 --> C23[Delta编码]

2.2 压缩算法对比

算法 压缩率 压缩速度 解压速度 适用场景
LZ4 极高 极高 实时数据
ZSTD 通用场景
Gzip 中等 中等 文件压缩
Snappy 极高 极高 内存压缩
Brotli 中等 Web传输

三、LZ4压缩算法

3.1 LZ4压缩实现

public class Lz4CompressionService
{
    public byte[] Compress(byte[] data)
    {
        var compressed = new byte[data.Length];
        
        var compressedSize = Lz4.Lz4Compress(data, compressed, data.Length, compressed.Length);
        
        return compressed.Take(compressedSize).ToArray();
    }
    
    public byte[] Decompress(byte[] compressedData, int originalSize)
    {
        var decompressed = new byte[originalSize];
        
        Lz4.Lz4Decompress(compressedData, decompressed, compressedData.Length, originalSize);
        
        return decompressed;
    }
    
    public async Task CompressAsync(byte[] data)
    {
        return await Task.Run(() => Compress(data));
    }
    
    public async Task DecompressAsync(byte[] compressedData, int originalSize)
    {
        return await Task.Run(() => Decompress(compressedData, originalSize));
    }
    
    public Stream CompressStream(Stream inputStream)
    {
        var memoryStream = new MemoryStream();
        
        using var lz4Stream = new Lz4Stream(memoryStream, CompressionMode.Compress);
        
        inputStream.CopyTo(lz4Stream);
        lz4Stream.Flush();
        
        memoryStream.Position = 0;
        
        return memoryStream;
    }
    
    public Stream DecompressStream(Stream compressedStream)
    {
        var memoryStream = new MemoryStream();
        
        using var lz4Stream = new Lz4Stream(compressedStream, CompressionMode.Decompress);
        
        lz4Stream.CopyTo(memoryStream);
        memoryStream.Position = 0;
        
        return memoryStream;
    }
}

四、ZSTD压缩算法

4.1 ZSTD压缩实现

public class ZstdCompressionService
{
    public byte[] Compress(byte[] data, int compressionLevel = 3)
    {
        var compressedSize = Zstd.ZstdCompressBound(data.Length);
        var compressed = new byte[compressedSize];
        
        var actualSize = Zstd.ZstdCompress(compressed, compressed.Length, data, data.Length, compressionLevel);
        
        return compressed.Take(actualSize).ToArray();
    }
    
    public byte[] Decompress(byte[] compressedData, int originalSize)
    {
        var decompressed = new byte[originalSize];
        
        Zstd.ZstdDecompress(decompressed, originalSize, compressedData, compressedData.Length);
        
        return decompressed;
    }
    
    public async Task CompressAsync(byte[] data, int compressionLevel = 3)
    {
        return await Task.Run(() => Compress(data, compressionLevel));
    }
    
    public async Task DecompressAsync(byte[] compressedData, int originalSize)
    {
        return await Task.Run(() => Decompress(compressedData, originalSize));
    }
    
    public byte[] CompressWithDictionary(byte[] data, byte[] dictionary, int compressionLevel = 3)
    {
        var compressedSize = Zstd.ZstdCompressBound(data.Length);
        var compressed = new byte[compressedSize];
        
        var actualSize = Zstd.ZstdCompressUsingDict(compressed, compressed.Length, 
            data, data.Length, dictionary, dictionary.Length, compressionLevel);
        
        return compressed.Take(actualSize).ToArray();
    }
    
    public byte[] DecompressWithDictionary(byte[] compressedData, byte[] dictionary, int originalSize)
    {
        var decompressed = new byte[originalSize];
        
        Zstd.ZstdDecompressUsingDict(decompressed, originalSize, compressedData, 
            compressedData.Length, dictionary, dictionary.Length);
        
        return decompressed;
    }
}

4.2 ZSTD字典训练

public class ZstdDictionaryService
{
    public byte[] TrainDictionary(List samples, int dictionarySize = 1000000)
    {
        var dictionary = new byte[dictionarySize];
        
        var totalSize = samples.Sum(s => s.Length);
        var concatenated = new byte[totalSize];
        var offset = 0;
        
        foreach (var sample in samples)
        {
            sample.CopyTo(concatenated, offset);
            offset += sample.Length;
        }
        
        Zstd.ZstdTrainFromBuffer(dictionary, dictionarySize, concatenated, totalSize);
        
        return dictionary;
    }
    
    public async Task TrainDictionaryAsync(List samples, int dictionarySize = 1000000)
    {
        return await Task.Run(() => TrainDictionary(samples, dictionarySize));
    }
}

五、列式压缩

5.1 列式压缩原理

graph TD A[行式存储] --> B[Row1: id=1, name='Alice', age=25] A --> C[Row2: id=2, name='Bob', age=30] D[列式存储] --> E[Column id: [1, 2]] D --> F[Column name: ['Alice', 'Bob']] D --> G[Column age: [25, 30]] E --> H[字典编码] H --> I[['Alice', 'Bob'] -> [1, 2]] G --> J[Delta编码] J --> K[[25, 30] -> [25, 5]] J --> L[行程编码] L --> M[[25, 25, 30] -> [(25, 2), (30, 1)]]

5.2 列式压缩实现

public class ColumnarCompressionService
{
    public byte[] DictionaryEncode(List values)
    {
        var dictionary = values.Distinct().ToList();
        var indices = values.Select(v => dictionary.IndexOf(v)).ToArray();
        
        var serialized = new MemoryStream();
        
        using var writer = new BinaryWriter(serialized);
        writer.Write(dictionary.Count);
        
        foreach (var word in dictionary)
        {
            writer.Write(word);
        }
        
        writer.Write(indices.Length);
        foreach (var index in indices)
        {
            writer.Write(index);
        }
        
        return serialized.ToArray();
    }
    
    public List DictionaryDecode(byte[] encoded)
    {
        var dictionary = new List();
        var indices = new List();
        
        using var reader = new BinaryReader(new MemoryStream(encoded));
        
        var dictionaryCount = reader.ReadInt32();
        for (int i = 0; i < dictionaryCount; i++)
        {
            dictionary.Add(reader.ReadString());
        }
        
        var indicesCount = reader.ReadInt32();
        for (int i = 0; i < indicesCount; i++)
        {
            indices.Add(reader.ReadInt32());
        }
        
        return indices.Select(i => dictionary[i]).ToList();
    }
    
    public byte[] DeltaEncode(List values)
    {
        var deltas = new List();
        
        if (values.Count > 0)
        {
            deltas.Add(values[0]);
            
            for (int i = 1; i < values.Count; i++)
            {
                deltas.Add(values[i] - values[i - 1]);
            }
        }
        
        var serialized = new MemoryStream();
        
        using var writer = new BinaryWriter(serialized);
        writer.Write(deltas.Count);
        foreach (var delta in deltas)
        {
            writer.Write(delta);
        }
        
        return serialized.ToArray();
    }
    
    public List DeltaDecode(byte[] encoded)
    {
        var deltas = new List();
        var values = new List();
        
        using var reader = new BinaryReader(new MemoryStream(encoded));
        
        var count = reader.ReadInt32();
        for (int i = 0; i < count; i++)
        {
            deltas.Add(reader.ReadInt32());
        }
        
        if (deltas.Count > 0)
        {
            values.Add(deltas[0]);
            
            for (int i = 1; i < deltas.Count; i++)
            {
                values.Add(values[i - 1] + deltas[i]);
            }
        }
        
        return values;
    }
    
    public byte[] RunLengthEncode(List values)
    {
        var runs = new List<(int Value, int Count)>();
        
        if (values.Count > 0)
        {
            var currentValue = values[0];
            var count = 1;
            
            for (int i = 1; i < values.Count; i++)
            {
                if (values[i] == currentValue)
                {
                    count++;
                }
                else
                {
                    runs.Add((currentValue, count));
                    currentValue = values[i];
                    count = 1;
                }
            }
            
            runs.Add((currentValue, count));
        }
        
        var serialized = new MemoryStream();
        
        using var writer = new BinaryWriter(serialized);
        writer.Write(runs.Count);
        foreach (var (value, runCount) in runs)
        {
            writer.Write(value);
            writer.Write(runCount);
        }
        
        return serialized.ToArray();
    }
}

六、数据压缩监控与优化

6.1 压缩监控

public class CompressionMonitor
{
    public CompressionMetrics MeasureCompression(byte[] data, CompressionAlgorithm algorithm)
    {
        var originalSize = data.Length;
        
        byte[] compressed;
        
        switch (algorithm)
        {
            case CompressionAlgorithm.Lz4:
                compressed = _lz4Service.Compress(data);
                break;
            case CompressionAlgorithm.Zstd:
                compressed = _zstdService.Compress(data);
                break;
            case CompressionAlgorithm.Gzip:
                compressed = CompressWithGzip(data);
                break;
            default:
                throw new ArgumentException("Unknown algorithm");
        }
        
        var compressedSize = compressed.Length;
        var compressionRatio = (double)compressedSize / originalSize;
        var savingsPercent = (1 - compressionRatio) * 100;
        
        return new CompressionMetrics
        {
            Algorithm = algorithm,
            OriginalSize = originalSize,
            CompressedSize = compressedSize,
            CompressionRatio = compressionRatio,
            SavingsPercent = savingsPercent
        };
    }
    
    private byte[] CompressWithGzip(byte[] data)
    {
        using var memoryStream = new MemoryStream();
        using var gzipStream = new GZipStream(memoryStream, CompressionLevel.Optimal);
        
        gzipStream.Write(data, 0, data.Length);
        gzipStream.Flush();
        
        return memoryStream.ToArray();
    }
}

public class CompressionMetrics
{
    public CompressionAlgorithm Algorithm { get; set; }
    public int OriginalSize { get; set; }
    public int CompressedSize { get; set; }
    public double CompressionRatio { get; set; }
    public double SavingsPercent { get; set; }
}

public enum CompressionAlgorithm { Lz4, Zstd, Gzip, Snappy, Brotli }

6.2 压缩优化策略

public class CompressionOptimizer
{
    public CompressionAlgorithm SelectBestAlgorithm(byte[] data)
    {
        var metrics = new List();
        
        foreach (var algorithm in Enum.GetValues())
        {
            metrics.Add(_compressionMonitor.MeasureCompression(data, algorithm));
        }
        
        return metrics.OrderBy(m => m.CompressionRatio).First().Algorithm;
    }
    
    public int SelectBestCompressionLevel(CompressionAlgorithm algorithm, CompressionPreference preference)
    {
        return (algorithm, preference) switch
        {
            (CompressionAlgorithm.Zstd, CompressionPreference.Speed) => 1,
            (CompressionAlgorithm.Zstd, CompressionPreference.Balance) => 3,
            (CompressionAlgorithm.Zstd, CompressionPreference.Compression) => 10,
            (CompressionAlgorithm.Gzip, CompressionPreference.Speed) => 1,
            (CompressionAlgorithm.Gzip, CompressionPreference.Balance) => 5,
            (CompressionAlgorithm.Gzip, CompressionPreference.Compression) => 9,
            _ => 3
        };
    }
    
    public async Task OptimizeCompressionAsync(byte[] data)
    {
        var algorithm = SelectBestAlgorithm(data);
        
        return algorithm switch
        {
            CompressionAlgorithm.Lz4 => _lz4Service.Compress(data),
            CompressionAlgorithm.Zstd => _zstdService.Compress(data),
            CompressionAlgorithm.Gzip => CompressWithGzip(data),
            _ => data
        };
    }
}

public enum CompressionPreference { Speed, Balance, Compression }

七、存储优化

7.1 存储优化策略

策略 描述 效果 适用场景
列式存储 按列存储数据 高压缩率 数据分析
分区存储 按条件分区 快速查询 时间序列
索引优化 合理创建索引 查询加速 通用场景
数据归档 归档历史数据 节省空间 冷数据
数据去重 去除重复数据 节省空间 日志存储

7.2 列式存储实现

public class ColumnarStorageService
{
    public ColumnarTable CreateColumnarTable(string tableName, List columns)
    {
        return new ColumnarTable
        {
            TableName = tableName,
            Columns = columns.ToDictionary(c => c.Name, c => c),
            ColumnData = new Dictionary()
        };
    }
    
    public void InsertData(ColumnarTable table, Dictionary row)
    {
        foreach (var (columnName, value) in row)
        {
            if (!table.Columns.ContainsKey(columnName))
            {
                continue;
            }
            
            var columnData = table.ColumnData.TryGetValue(columnName, out var data) 
                ? data 
                : Array.Empty();
            
            var encodedValue = EncodeValue(value, table.Columns[columnName].DataType);
            table.ColumnData[columnName] = columnData.Concat(encodedValue).ToArray();
        }
    }
    
    public List> QueryData(ColumnarTable table, List columns)
    {
        var result = new List>();
        
        var rowCount = table.ColumnData.Any() ? GetRowCount(table) : 0;
        
        for (int i = 0; i < rowCount; i++)
        {
            var row = new Dictionary();
            
            foreach (var columnName in columns)
            {
                if (table.ColumnData.TryGetValue(columnName, out var data))
                {
                    row[columnName] = DecodeValue(data, table.Columns[columnName].DataType, i);
                }
            }
            
            result.Add(row);
        }
        
        return result;
    }
    
    private int GetRowCount(ColumnarTable table)
    {
        var firstColumn = table.ColumnData.First();
        
        return firstColumn.Value.Length / GetColumnSize(table.Columns[firstColumn.Key].DataType);
    }
    
    private int GetColumnSize(DataType dataType)
    {
        return dataType switch
        {
            DataType.Int32 => 4,
            DataType.Int64 => 8,
            DataType.Float => 4,
            DataType.Double => 8,
            _ => 1
        };
    }
}

public class ColumnarTable
{
    public string TableName { get; set; }
    public Dictionary Columns { get; set; }
    public Dictionary ColumnData { get; set; }
}

public class ColumnDefinition
{
    public string Name { get; set; }
    public DataType DataType { get; set; }
    public bool Compressed { get; set; }
}

public enum DataType { Int32, Int64, Float, Double, String, DateTime }

八、数据压缩最佳实践

8.1 压缩算法选型

场景 推荐算法 原因
实时数据 LZ4 速度最快
通用场景 ZSTD 平衡
文件存储 Gzip 兼容性
内存压缩 Snappy 低CPU
Web传输 Brotli 高压缩率

8.2 压缩实施建议

  • 选择合适的压缩算法
  • 根据数据类型选择压缩级别
  • 使用列式存储优化分析场景
  • 定期监控压缩效果
  • 考虑压缩/解压开销

8.3 存储优化最佳实践

public class StorageOptimizationService
{
    public async Task OptimizeStorageAsync(StorageOptimizationOptions options)
    {
        foreach (var table in options.Tables)
        {
            await CompressTableAsync(table);
            await PartitionTableAsync(table);
            await ArchiveOldDataAsync(table);
        }
    }
    
    private async Task CompressTableAsync(string tableName)
    {
        var data = await _database.GetTableDataAsync(tableName);
        var compressed = await _compressionService.CompressAsync(data);
        await _database.SaveCompressedDataAsync(tableName, compressed);
    }
    
    private async Task PartitionTableAsync(string tableName)
    {
        await _database.PartitionTableAsync(tableName, PartitionType.Monthly);
    }
    
    private async Task ArchiveOldDataAsync(string tableName)
    {
        await _database.ArchiveDataAsync(tableName, DateTime.UtcNow.AddYears(-1));
    }
}

public class StorageOptimizationOptions
{
    public List Tables { get; set; }
    public CompressionAlgorithm CompressionAlgorithm { get; set; }
    public PartitionType PartitionType { get; set; }
}

public enum PartitionType { Daily, Weekly, Monthly, Yearly }

九、总结

数据压缩是减少存储空间和网络传输开销的关键技术。LZ4适合实时数据,ZSTD适合通用场景,Gzip适合文件存储。通过合理选择压缩算法、实现列式存储、做好监控和优化,能够构建高效的数据存储系统。定期优化存储配置、归档历史数据、实施数据去重,能够显著减少存储成本。