一、列式存储概述
列式存储是一种按列存储数据的方式,与传统的行式存储相对。列式存储在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查询效率,数据压缩能够显著减少存储空间占用。合理选择压缩算法、使用字典编码和位图索引等优化技术,能够构建高效的数据存储系统,为数据分析提供有力支持。