📖 数据密集型设计

数据压缩与列式存储

深入探讨数据压缩技术与列式存储优化

一、列式存储概述

列式存储是一种按列存储数据的方式,与传统的行式存储相对。列式存储在OLAP场景下具有显著的性能优势,能够大幅提高查询效率和数据压缩率。

二、行式存储 vs 列式存储

2.1 存储方式对比

graph TD subgraph 行式存储 A[Row 1] --> A1[id=1] A --> A2[name=Tom] A --> A3[age=25] A --> A4[score=90] B[Row 2] --> B1[id=2] B --> B2[name=Jerry] B --> B3[age=23] B --> B4[score=85] end subgraph 列式存储 C[Column id] --> C1[1] C --> C2[2] D[Column name] --> D1[Tom] D --> D2[Jerry] E[Column age] --> E1[25] E --> E2[23] F[Column score] --> F1[90] F --> F2[85] end

2.2 存储方式对比表

特性 行式存储 列式存储
存储方式 按行连续存储 按列连续存储
查询方式 读取整行数据 只读取需要的列
压缩率 较低 较高
写入性能 较高 较低
适用场景 OLTP、事务处理 OLAP、数据分析

三、列式存储实现原理

3.1 列式存储架构

graph TD A[列式存储引擎] --> B[数据写入] A --> C[数据读取] B --> D[数据缓冲] D --> E[列压缩] E --> F[列存储] C --> G[列解压] G --> H[列投影] H --> I[查询结果]

3.2 Parquet列式存储

public class ParquetStorageService
{
    public async Task WriteParquetAsync<T>(string filePath, List<T> data)
    {
        using var stream = File.OpenWrite(filePath);
        
        using var parquetWriter = new ParquetWriter(
            ParquetSchemaBuilder.Build<T>(),
            stream
        );
        
        parquetWriter.CompressionMethod = CompressionMethod.Snappy;
        
        using var rowGroupWriter = parquetWriter.CreateRowGroup();
        await rowGroupWriter.WriteAsync(data);
    }
    
    public async Task<List<T>> ReadParquetAsync<T>(string filePath)
    {
        using var stream = File.OpenRead(filePath);
        
        using var parquetReader = new ParquetReader(stream);
        
        var result = new List<T>();
        
        for (int i = 0; i < parquetReader.RowGroupCount; i++)
        {
            using var rowGroupReader = parquetReader.OpenRowGroupReader(i);
            var rowGroupData = await rowGroupReader.ReadAsync<T>();
            result.AddRange(rowGroupData);
        }
        
        return result;
    }
    
    public async Task<List<T>> ReadParquetWithProjectionAsync<T>(
        string filePath, 
        List<string> columnsToRead)
    {
        using var stream = File.OpenRead(filePath);
        
        using var parquetReader = new ParquetReader(stream);
        
        var schema = parquetReader.Schema;
        var projection = schema.GetDataFields()
            .Where(f => columnsToRead.Contains(f.Name))
            .ToList();
        
        var result = new List<T>();
        
        for (int i = 0; i < parquetReader.RowGroupCount; i++)
        {
            using var rowGroupReader = parquetReader.OpenRowGroupReader(i);
            var rowGroupData = await rowGroupReader.ReadAsync<T>(projection);
            result.AddRange(rowGroupData);
        }
        
        return result;
    }
}

3.3 ORC列式存储

public class OrcStorageService
{
    public async Task WriteOrcAsync(string filePath, DataTable data)
    {
        using var stream = File.OpenWrite(filePath);
        
        var writer = OrcFile.CreateWriter(stream, 
            OrcFile.CreateOptions(data.Schema));
        
        writer.AddRowGroup(data);
        writer.Close();
    }
    
    public async Task<DataTable> ReadOrcAsync(string filePath)
    {
        using var stream = File.OpenRead(filePath);
        
        var reader = OrcFile.CreateReader(stream, OrcFile.ReadOptions.Default);
        
        return reader.RowIterator().Select(row => row).CopyToDataTable();
    }
}

四、数据压缩算法

4.1 压缩算法对比

算法 压缩率 压缩速度 解压速度 适用场景
LZ4 极快 极快 实时数据
ZSTD 综合场景
Snappy 大数据处理
GZIP 较高 归档存储
Brotli 极高 静态资源

4.2 LZ4压缩实现

public class Lz4CompressionService
{
    public byte[] Compress(byte[] data)
    {
        var maxCompressedLength = LZ4Codec.MaximumOutputLength(data.Length);
        var compressed = new byte[maxCompressedLength];
        
        var compressedLength = LZ4Codec.Encode(
            data, 0, data.Length,
            compressed, 0, maxCompressedLength
        );
        
        Array.Resize(ref compressed, compressedLength);
        return compressed;
    }
    
    public byte[] Decompress(byte[] compressedData, int originalLength)
    {
        var decompressed = new byte[originalLength];
        
        LZ4Codec.Decode(
            compressedData, 0, compressedData.Length,
            decompressed, 0, originalLength
        );
        
        return decompressed;
    }
    
    public async Task<byte[]> CompressAsync(byte[] data)
    {
        return await Task.Run(() => Compress(data));
    }
    
    public async Task<byte[]> DecompressAsync(byte[] compressedData, int originalLength)
    {
        return await Task.Run(() => Decompress(compressedData, originalLength));
    }
}

4.3 ZSTD压缩实现

public class ZstdCompressionService
{
    private readonly int _compressionLevel;
    
    public ZstdCompressionService(int compressionLevel = 3)
    {
        _compressionLevel = compressionLevel;
    }
    
    public byte[] Compress(byte[] data)
    {
        var maxCompressedLength = ZstdNet.Compressor.MaxCompressedSize(data.Length);
        var compressed = new byte[maxCompressedLength];
        
        var compressor = new ZstdNet.Compressor(_compressionLevel);
        var compressedLength = compressor.Wrap(data, compressed);
        
        Array.Resize(ref compressed, compressedLength);
        return compressed;
    }
    
    public byte[] Decompress(byte[] compressedData)
    {
        var decompressor = new ZstdNet.Decompressor();
        return decompressor.Unwrap(compressedData);
    }
    
    public async Task CompressFileAsync(string inputPath, string outputPath)
    {
        using var inputStream = File.OpenRead(inputPath);
        using var outputStream = File.OpenWrite(outputPath);
        using var compressionStream = new ZstdNet.CompressionStream(outputStream, _compressionLevel);
        
        await inputStream.CopyToAsync(compressionStream);
    }
    
    public async Task DecompressFileAsync(string inputPath, string outputPath)
    {
        using var inputStream = File.OpenRead(inputPath);
        using var decompressionStream = new ZstdNet.DecompressionStream(inputStream);
        using var outputStream = File.OpenWrite(outputPath);
        
        await decompressionStream.CopyToAsync(outputStream);
    }
}

五、列式存储优化技术

5.1 字典编码

字典编码将重复值替换为索引:

graph TD A[原始数据] --> A1[北京,上海,北京,广州,上海] A1 --> B[字典构建] B --> C[字典: 0=北京, 1=上海, 2=广州] C --> D[编码后] D --> E[0, 1, 0, 2, 1]

5.2 字典编码实现

public class DictionaryEncodingService
{
    public DictionaryEncodingResult Encode(List<string> values)
    {
        var uniqueValues = values.Distinct().ToList();
        var dictionary = uniqueValues.Select((value, index) => new { value, index })
            .ToDictionary(x => x.value, x => x.index);
        
        var encoded = values.Select(v => dictionary[v]).ToList();
        
        return new DictionaryEncodingResult
        {
            Dictionary = uniqueValues,
            EncodedData = encoded
        };
    }
    
    public List<string> Decode(DictionaryEncodingResult result)
    {
        return result.EncodedData.Select(index => result.Dictionary[index]).ToList();
    }
    
    public double GetCompressionRatio(List<string> original)
    {
        var result = Encode(original);
        
        var originalSize = original.Sum(v => v.Length) * sizeof(char);
        var encodedSize = result.EncodedData.Count * sizeof(int) + 
                          result.Dictionary.Sum(v => v.Length) * sizeof(char);
        
        return (double)originalSize / encodedSize;
    }
}

public class DictionaryEncodingResult
{
    public List<string> Dictionary { get; set; }
    public List<int> EncodedData { get; set; }
}

5.3 Run-Length Encoding (RLE)

public class RleEncodingService
{
    public List<(T Value, int Count)> Encode<T>(List<T> values) where T : IEquatable<T>
    {
        if (values == null || values.Count == 0)
            return new List<(T, int)>();
        
        var result = new List<(T, int)>();
        var currentValue = values[0];
        var count = 1;
        
        for (int i = 1; i < values.Count; i++)
        {
            if (values[i].Equals(currentValue))
            {
                count++;
            }
            else
            {
                result.Add((currentValue, count));
                currentValue = values[i];
                count = 1;
            }
        }
        
        result.Add((currentValue, count));
        return result;
    }
    
    public List<T> Decode<T>(List<(T Value, int Count)> encoded)
    {
        var result = new List<T>();
        
        foreach (var (value, count) in encoded)
        {
            result.AddRange(Enumerable.Repeat(value, count));
        }
        
        return result;
    }
}

5.4 位图索引

public class BitmapIndexService
{
    public Dictionary<object, BitArray> CreateBitmapIndex<T>(List<T> data, Func<T, object> selector)
    {
        var distinctValues = data.Select(selector).Distinct().ToList();
        var index = new Dictionary<object, BitArray>();
        
        foreach (var value in distinctValues)
        {
            var bitmap = new BitArray(data.Count);
            
            for (int i = 0; i < data.Count; i++)
            {
                if (selector(data[i]).Equals(value))
                {
                    bitmap.Set(i, true);
                }
            }
            
            index[value] = bitmap;
        }
        
        return index;
    }
    
    public List<int> QueryBitmapIndex(Dictionary<object, BitArray> index, object value)
    {
        if (!index.TryGetValue(value, out var bitmap))
            return new List<int>();
        
        var result = new List<int>();
        
        for (int i = 0; i < bitmap.Length; i++)
        {
            if (bitmap.Get(i))
            {
                result.Add(i);
            }
        }
        
        return result;
    }
}

六、压缩策略选择

6.1 压缩策略选择流程

flowchart TD A[选择压缩策略] --> B{数据类型} B -->|文本数据| C[ZSTD/Brotli] B -->|数值数据| D[LZ4/ZSTD] B -->|重复数据| E[字典编码+RLE] A --> F{性能要求} F -->|实时处理| G[LZ4/Snappy] F -->|批量处理| H[ZSTD] F -->|归档存储| I[GZIP/Brotli] A --> G{压缩率要求} G -->|高压缩率| J[ZSTD/Brotli] G -->|平衡| K[LZ4/Snappy]

6.2 动态压缩策略

public class DynamicCompressionStrategy
{
    public ICompressionService SelectCompressionService(DataCharacteristics characteristics)
    {
        if (characteristics.IsRealTime)
            return new Lz4CompressionService();
        
        if (characteristics.CompressionRatioPriority)
            return new ZstdCompressionService(9);
        
        if (characteristics.HasHighCardinality)
            return new ZstdCompressionService(3);
        
        if (characteristics.HasManyDuplicates)
            return new DictionaryEncodingService().CombineWith(new Lz4CompressionService());
        
        return new ZstdCompressionService(3);
    }
}

public class DataCharacteristics
{
    public bool IsRealTime { get; set; }
    public bool CompressionRatioPriority { get; set; }
    public bool HasHighCardinality { get; set; }
    public bool HasManyDuplicates { get; set; }
    public DataCategory Category { get; set; }
}

public enum DataCategory { Text, Numeric, Binary, Log }
}

七、列式存储最佳实践

7.1 选择合适的列式格式

根据场景选择Parquet、ORC或其他列式格式。

7.2 合理设置压缩级别

在压缩率和性能之间取得平衡。

7.3 使用字典编码优化

对低基数列使用字典编码。

7.4 分区和分桶

public class ColumnarStorageOptimizer
{
    public async Task OptimizeStorage(string tableName, StorageConfig config)
    {
        await ApplyPartitioning(tableName, config.PartitionConfig);
        await ApplyBucketing(tableName, config.BucketConfig);
        await ApplyCompression(tableName, config.CompressionConfig);
        await ApplyIndexing(tableName, config.IndexConfig);
    }
    
    private async Task ApplyPartitioning(string tableName, PartitionConfig config)
    {
        var sql = $"ALTER TABLE {tableName} PARTITION BY {config.Column}";
        await _dbContext.Database.ExecuteSqlRawAsync(sql);
    }
    
    private async Task ApplyCompression(string tableName, CompressionConfig config)
    {
        var sql = $"ALTER TABLE {tableName} SET COMPRESSION '{config.Algorithm}'";
        await _dbContext.Database.ExecuteSqlRawAsync(sql);
    }
}

八、列式存储性能优化

8.1 查询优化

-- 好的查询:只查询需要的列
SELECT user_id, order_date, total_amount FROM orders WHERE order_date >= '2024-01-01';

-- 好的查询:使用谓词下推
SELECT * FROM orders WHERE category = 'electronics' AND total_amount > 100;

-- 差的查询:SELECT *
SELECT * FROM orders;

-- 差的查询:没有谓词
SELECT user_id FROM orders;

8.2 数据统计优化

public class ColumnarStatisticsService
{
    public async Task UpdateStatistics(string tableName)
    {
        await _dbContext.Database.ExecuteSqlRawAsync(
            $"ANALYZE TABLE {tableName}"
        );
    }
    
    public async Task<ColumnStatistics> GetColumnStatistics(string tableName, string columnName)
    {
        var stats = await _dbContext.QuerySingleAsync<ColumnStatistics>(
            $"SELECT COUNT(DISTINCT {columnName}) as DistinctCount, " +
            $"MIN({columnName}) as MinValue, " +
            $"MAX({columnName}) as MaxValue " +
            $"FROM {tableName}"
        );
        
        return stats;
    }
}

九、列式存储监控

9.1 监控指标

public class ColumnarStorageMetrics
{
    public string TableName { get; set; }
    public long OriginalSize { get; set; }
    public long CompressedSize { get; set; }
    public double CompressionRatio { get; set; }
    public double QueryTimeMs { get; set; }
    public int ColumnCount { get; set; }
}

public class ColumnarStorageMonitor
{
    public async Task<ColumnarStorageMetrics> GetMetrics(string tableName)
    {
        var metrics = new ColumnarStorageMetrics { TableName = tableName };
        
        var sizeInfo = await _dbContext.QuerySingleAsync<SizeInfo>(
            $"SELECT data_length, index_length FROM information_schema.tables WHERE table_name = '{tableName}'"
        );
        
        metrics.OriginalSize = sizeInfo.DataLength;
        metrics.CompressedSize = await GetCompressedSize(tableName);
        metrics.CompressionRatio = (double)metrics.OriginalSize / metrics.CompressedSize;
        
        return metrics;
    }
}

十、总结

列式存储和数据压缩是构建高性能数据仓库的关键技术。列式存储能够大幅提高OLAP查询效率,数据压缩能够显著减少存储空间占用。合理选择压缩算法、使用字典编码和位图索引等优化技术,能够构建高效的数据存储系统,为数据分析提供有力支持。