📖 数据密集型设计

数据压缩与列式存储

深入探讨数据压缩技术与存储优化策略

一、数据压缩概述

数据压缩是通过减少数据占用空间来优化存储和传输的技术。在数据密集型应用中,数据压缩能够显著降低存储成本、提高IO效率。

二、压缩算法分类

2.1 压缩算法对比

算法 压缩率 压缩速度 解压速度 适用场景
LZ4 中等 极快 极快 实时数据、缓存
ZSTD 极快 通用场景
Gzip 中等 静态数据
Snappy 中等 极快 极快 大数据处理
Brotli 极高 中等 Web传输

2.2 压缩算法选择

graph TD A[选择压缩算法] --> B{压缩率优先} B -->|是| C[Gzip/Brotli] B -->|否| D{速度优先} D -->|是| E[LZ4/Snappy] D -->|否| F[ZSTD] C --> C1[静态数据存储] E --> E1[实时数据处理] F --> F1[通用场景]

三、LZ4压缩算法

3.1 LZ4原理

LZ4是一种基于LZ77的快速压缩算法,采用哈希表加速匹配查找:

flowchart TD A[输入数据] --> B[滑动窗口] B --> C[哈希表] C --> D{查找匹配} D -->|找到| E[编码为引用] D -->|未找到| F[编码为字面量] E --> G[输出压缩数据] F --> G

3.2 LZ4实现

public class Lz4Compressor
{
    public byte[] Compress(byte[] input)
    {
        using (var ms = new MemoryStream())
        using (var lz4Stream = new LZ4.LZ4Stream(ms, System.IO.Compression.CompressionMode.Compress))
        {
            lz4Stream.Write(input, 0, input.Length);
            lz4Stream.Flush();
            return ms.ToArray();
        }
    }
    
    public byte[] Decompress(byte[] compressed)
    {
        using (var ms = new MemoryStream(compressed))
        using (var lz4Stream = new LZ4.LZ4Stream(ms, System.IO.Compression.CompressionMode.Decompress))
        using (var outputMs = new MemoryStream())
        {
            lz4Stream.CopyTo(outputMs);
            return outputMs.ToArray();
        }
    }
}

3.3 LZ4特性

  • 压缩速度:~400-800 MB/s
  • 解压速度:~1 GB/s+
  • 压缩率:通常2-4倍
  • 内存占用:低

四、ZSTD压缩算法

4.1 ZSTD原理

ZSTD(Zstandard)是Facebook开发的新一代压缩算法,结合了LZ77和FSE(Finite State Entropy):

flowchart TD A[输入数据] --> B[LZ77匹配] B --> C[字面量序列] B --> D[匹配序列] C --> E[FSE编码] D --> F[FSE编码] E --> G[输出压缩数据] F --> G

4.2 ZSTD实现

public class ZstdCompressor
{
    public byte[] Compress(byte[] input, int level = 3)
    {
        var compressed = new byte[ZstdNet.Zstd.GetMaxCompressedSize(input.Length)];
        var compressedSize = ZstdNet.Zstd.Compress(input, compressed, level);
        return compressed.Take(compressedSize).ToArray();
    }
    
    public byte[] Decompress(byte[] compressed)
    {
        var decompressed = new byte[1024 * 1024]; // 预估大小
        var decompressedSize = ZstdNet.Zstd.Decompress(compressed, decompressed);
        return decompressed.Take(decompressedSize).ToArray();
    }
}

4.3 ZSTD特性

  • 压缩速度:~100-500 MB/s(取决于级别)
  • 解压速度:~500-1500 MB/s
  • 压缩率:通常3-5倍
  • 可配置压缩级别:1-22

五、列式存储

5.1 列式存储vs行式存储

graph TD A[行式存储] --> B[Row 1: col1, col2, col3] A --> C[Row 2: col1, col2, col3] A --> D[Row 3: col1, col2, col3] E[列式存储] --> F[Column 1: row1, row2, row3] E --> G[Column 2: row1, row2, row3] E --> H[Column 3: row1, row2, row3]

5.2 列式存储优势

特性 行式存储 列式存储
查询类型 OLTP OLAP
压缩率
列查询效率
写入效率

5.3 列式存储实现

public class ColumnarStorage
{
    private readonly Dictionary<string, List<object>> _columns = new Dictionary<string, List<object>>();
    
    public void AddColumn(string name, List<object> values)
    {
        _columns[name] = values;
    }
    
    public List<object> GetColumn(string name)
    {
        return _columns.TryGetValue(name, out var column) ? column : null;
    }
    
    public double GetColumnCompressionRatio(string name)
    {
        var column = GetColumn(name);
        if (column == null)
            return 0;
        
        var rawSize = column.Sum(v => GetValueSize(v));
        var compressed = CompressColumn(column);
        var compressedSize = compressed.Length;
        
        return (double)rawSize / compressedSize;
    }
    
    private byte[] CompressColumn(List<object> column)
    {
        var data = JsonSerializer.Serialize(column);
        var compressor = new ZstdCompressor();
        return compressor.Compress(Encoding.UTF8.GetBytes(data));
    }
    
    private int GetValueSize(object value)
    {
        return Encoding.UTF8.GetByteCount(value.ToString());
    }
}

六、字典编码

6.1 字典编码原理

字典编码将重复出现的值映射为整数,减少存储空间:

graph TD A[原始数据] --> B[提取唯一值] B --> C[构建字典] C --> D[值: 索引] D --> D1[北京: 0] D --> D2[上海: 1] D --> D3[广州: 2] A --> E[替换为索引] E --> F[压缩后数据]

6.2 字典编码实现

public class DictionaryEncoder
{
    private Dictionary<string, int> _valueToIndex = new Dictionary<string, int>();
    private List<string> _indexToValue = new List<string>();
    
    public int Encode(string value)
    {
        if (!_valueToIndex.TryGetValue(value, out var index))
        {
            index = _indexToValue.Count;
            _valueToIndex[value] = index;
            _indexToValue.Add(value);
        }
        return index;
    }
    
    public string Decode(int index)
    {
        return _indexToValue[index];
    }
    
    public byte[] GetDictionary()
    {
        return Encoding.UTF8.GetBytes(JsonSerializer.Serialize(_indexToValue));
    }
    
    public void LoadDictionary(byte[] dictionary)
    {
        var values = JsonSerializer.Deserialize<List<string>>(dictionary);
        _indexToValue = values;
        _valueToIndex = values.Select((v, i) => (v, i)).ToDictionary(x => x.v, x => x.i);
    }
}

6.3 字典编码应用

// 示例:压缩城市数据
var cities = new[] { "北京", "上海", "广州", "北京", "上海", "深圳" };
var encoder = new DictionaryEncoder();

// 编码
var encoded = cities.Select(c => encoder.Encode(c)).ToArray(); // [0, 1, 2, 0, 1, 3]

// 解码
var decoded = encoded.Select(i => encoder.Decode(i)).ToArray(); // ["北京", "上海", "广州", "北京", "上海", "深圳"]

// 字典大小
var dictionarySize = encoder.GetDictionary().Length; // 包含所有唯一值

七、数据压缩应用场景

7.1 数据库压缩

-- MySQL压缩表
CREATE TABLE compressed_table (
    id INT,
    data TEXT
) ENGINE=InnoDB ROW_FORMAT=COMPRESSED;

-- PostgreSQL压缩
CREATE TABLE compressed_table (
    id INT,
    data TEXT
) WITH (COMPRESSION = 'zstd');

-- SQLite压缩
PRAGMA page_size = 4096;
PRAGMA journal_mode = WAL;
PRAGMA synchronous = NORMAL;

7.2 缓存压缩

// Redis压缩
public class CompressedRedisCache
{
    private readonly IDistributedCache _cache;
    private readonly ICompressor _compressor;
    
    public async Task SetAsync(string key, object value)
    {
        var data = JsonSerializer.Serialize(value);
        var compressed = _compressor.Compress(Encoding.UTF8.GetBytes(data));
        await _cache.SetAsync(key, compressed);
    }
    
    public async Task<T> GetAsync<T>(string key)
    {
        var compressed = await _cache.GetAsync(key);
        if (compressed == null)
            return default;
        
        var data = _compressor.Decompress(compressed);
        return JsonSerializer.Deserialize<T>(data);
    }
}

7.3 消息队列压缩

// Kafka压缩配置
properties.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "zstd");

// RabbitMQ压缩
var compressedBody = compressor.Compress(messageBody);
channel.BasicPublish(exchange, routingKey, basicProperties, compressedBody);

// 消费者解压
var decompressedBody = compressor.Decompress(delivery.Body);

7.4 网络传输压缩

// HTTP压缩 (ASP.NET Core)
app.UseResponseCompression();

// 配置压缩中间件
services.AddResponseCompression(options =>
{
    options.Providers.Add<BrotliCompressionProvider>();
    options.Providers.Add<GzipCompressionProvider>();
    options.MimeTypes = ResponseCompressionDefaults.MimeTypes.Concat(new[] { "application/json" });
});

// 自定义压缩
public class CompressionMiddleware
{
    public async Task InvokeAsync(HttpContext context)
    {
        var originalBody = context.Response.Body;
        
        using (var compressedStream = new MemoryStream())
        {
            context.Response.Body = compressedStream;
            await _next(context);
            
            compressedStream.Seek(0, SeekOrigin.Begin);
            var compressed = compressor.Compress(compressedStream.ToArray());
            
            context.Response.Headers["Content-Encoding"] = "zstd";
            context.Response.ContentLength = compressed.Length;
            await context.Response.Body.WriteAsync(compressed);
        }
        
        context.Response.Body = originalBody;
    }
}

八、压缩策略选择

8.1 实时数据场景

实时数据处理需要快速压缩和解压,推荐使用LZ4或Snappy:

flowchart TD A[实时数据] --> B[数据采集] B --> C[LZ4压缩] C --> D[传输] D --> E[LZ4解压] E --> F[实时处理]

8.2 存储场景

存储场景追求高压缩率,推荐使用ZSTD或Gzip:

flowchart TD A[数据写入] --> B[ZSTD压缩] B --> C[存储] D[数据读取] --> E[ZSTD解压] E --> F[数据使用]

8.3 网络传输场景

网络传输场景需要平衡压缩率和速度,推荐使用Brotli或Gzip:

flowchart TD A[服务器] --> B[Brotli压缩] B --> C[HTTP传输] C --> D[浏览器解压] D --> E[渲染页面]

九、压缩监控

9.1 监控指标

public class CompressionMetrics
{
    public string Algorithm { get; set; }
    public double CompressionRatio { get; set; }
    public long CompressionTimeMs { get; set; }
    public long DecompressionTimeMs { get; set; }
    public long OriginalSize { get; set; }
    public long CompressedSize { get; set; }
}

public class CompressionMonitor
{
    public CompressionMetrics MeasureCompression(byte[] data, ICompressor compressor)
    {
        var sw = Stopwatch.StartNew();
        var compressed = compressor.Compress(data);
        var compressionTime = sw.ElapsedMilliseconds;
        
        sw.Restart();
        var decompressed = compressor.Decompress(compressed);
        var decompressionTime = sw.ElapsedMilliseconds;
        
        return new CompressionMetrics
        {
            Algorithm = compressor.GetType().Name,
            CompressionRatio = (double)data.Length / compressed.Length,
            CompressionTimeMs = compressionTime,
            DecompressionTimeMs = decompressionTime,
            OriginalSize = data.Length,
            CompressedSize = compressed.Length
        };
    }
}

十、总结

数据压缩是数据密集型应用中重要的性能优化手段。LZ4适合实时场景,ZSTD适合通用场景,列式存储和字典编码能够进一步提升压缩率。根据业务场景选择合适的压缩策略,能够显著降低存储成本和提高系统性能。