📖 数据密集型设计

数据压缩与列式存储

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

一、数据压缩概述

数据压缩是通过减少数据占用空间来优化存储和传输的技术。在数据密集型应用中,数据压缩能够显著降低存储成本、提高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适合通用场景,列式存储和字典编码能够进一步提升压缩率。根据业务场景选择合适的压缩策略,能够显著降低存储成本和提高系统性能。