Grafana Tempo 中的 zstd 压缩实战:klauspost/compress/zstd 库用法、并发模型与源码级解析
【免费下载链接】tempoGrafana Tempo is a high volume, minimal dependency distributed tracing backend.项目地址: https://gitcode.com/GitHub_Trending/tempo1/tempo
zstd(Zstandard)是一种兼顾高压缩比与实时解压速度的压缩算法,其 Go 实现github.com/klauspost/compress/zstd以纯 Go 编写、面向高吞吐场景做了深度优化,是 Grafana Tempo 分布式追踪后端在块存储、租户索引与 Kafka 消息压缩等多个环节实际依赖的核心压缩库。本文以该库的官方 README 为骨架,结合 Tempo 仓库中的真实调用源码,系统讲解其安装方式、压缩/解压 API、并发模型、字典压缩与性能基准,帮助你理解并复用这套高性能压缩方案。
一、zstd 与 klauspost/compress 概览
Zstandard 是 Facebook 开源的一种实时压缩算法,核心特点是在极广的压缩比/速度权衡区间内可调,且配套的解码器速度极快。github.com/klauspost/compress/zstd包提供了对 Zstandard 内容的压缩(Compressor)与解压(Decompressor)能力,目前对压缩端侧重速度优化。
该包具有以下关键特性:
- 纯 Go 实现:可用
noasm和nounsafe构建标签关闭汇编优化与 unsafe 特性; - 面向 64 位处理器深度优化:文档明确指出该包针对 64 位处理器重度优化,在 32 位处理器上会明显变慢;
- 开源许可:以 Go 标准开源许可证发布。
在 Tempo 仓库中,该库被直接用于三个核心场景(详见后文「六、Tempo 中的真实调用链」),分别是块后端的 zstd 编解码器、基于 protobuf 的租户索引序列化,以及 Kafka 生产者的消息压缩。
二、安装与引入
标准 Go 模块安装方式:
go get -u github.com/klauspost/compress包位于github.com/klauspost/compress/zstd,在代码中引入:
import "github.com/klauspost/compress/zstd"Tempo 仓库通过go.mod依赖该库,并在vendor/github.com/klauspost/compress/zstd/下固化了对应版本源码,确保构建可复现。
三、压缩器(Compressor)详解
3.1 稳定性状态与压缩等级
压缩器当前状态为STABLE:已有大量不同类型的内容被测试,并被多个项目活跃使用,所有更新都会经过模糊测试(fuzz-tested)。文档同时提醒:仍可能存在特定数据类型/大小/参数组合下的边界情况,因此生产使用前建议自行测试。
当前实现了三个速度档位的压缩器,其压缩比与原生 zstd 等级的大致对应关系如下:
| 库内档位 | 对应 zstd 等级 | 说明 |
|---|---|---|
| Fastest(最快) | 约 level 1 | 追求速度 |
| Default(默认) | 约 level 3(zstd 默认档) | 速度与压缩比均衡 |
| Better(更好) | 约 level 7 | 更高压缩比 |
| Best(最佳) | 约 level 11 | 极致压缩比,速度最慢 |
与 Go 标准库 deflate/gzip 相比,该库在最快模式下通常快约 2 倍;压缩比大致相当于 stdlib level 3,但通常快约 3 倍。
3.2 流式压缩的基本用法
Encoder有两种使用方式:通过io.WriteCloser接口进行流式压缩,或通过EncodeAll处理多个独立的小任务。文档建议小块数据优先使用EncodeAll。使用NewWriter创建实例,可同时支持两种方式。
默认选项下创建压缩器并流式压缩:
// Compress input to output. func Compress(in io.Reader, out io.Writer) error { enc, err := zstd.NewWriter(out) if err != nil { return err } _, err = io.Copy(enc, in) if err != nil { enc.Close() return err } return enc.Close() }写入数据到enc即可完成编码,调用Close()时输出内容才全部写完。即使编码失败也应调用Close(),以释放可能占用的资源。
3.3 复用 Writer 与并发控制
大块编码场景下,尽量复用 writer:使用Reset(io.Writer)切换输出目标,可复用全部内部资源、避免浪费性分配。
默认情况下,流式编码采用"轻量"并发——最多 2 个 goroutine 并行处理流的一部分。若希望关闭异步 goroutine、按块完成即阻塞写入,可使用WithEncoderConcurrency(1)。
3.4 并行流压缩(Parallel Stream Compression)
追求大流量最大吞吐时,组合使用WithConcurrentBlocks(true)与WithEncoderConcurrency(n)(n 为想使用的 CPU 核数),将输入切分为大段任务由多个 goroutine 并行压缩,类似 C zstd 库的多线程压缩:
enc, err := zstd.NewWriter(out, zstd.WithEncoderLevel(zstd.SpeedDefault), zstd.WithEncoderConcurrency(runtime.GOMAXPROCS(0)), zstd.WithConcurrentBlocks(true), )工作方式与注意点:
- 每个非首个任务会从前一个任务获得重叠前缀作为匹配上下文,因此压缩比仅受轻微影响;
- 输出按顺序刷新,最终生成合法的单帧 zstd 流;
- 注意:并行流压缩与字典编码不兼容;
Flush()会派发当前的部分任务,对延迟敏感的场景可强制输出;EncodeAll不受影响——它通过编码器池使用自身的并发。
作者在 1.8GB GOB 流(AMD Ryzen 9 9950X)上的基准数据:
| Level | 1 thread | 4 threads | 16 threads | 1T ratio | 16T ratio |
|---|---|---|---|---|---|
| fastest | 783 MB/s | 2950 MB/s (3.8×) | 6939 MB/s (8.9×) | 12.24% | 12.26% |
| default | 728 MB/s | 2533 MB/s (3.5×) | 5340 MB/s (7.3×) | 10.67% | 10.68% |
| better | 434 MB/s | 1105 MB/s (2.5×) | 2206 MB/s (5.1×) | 9.14% | 9.21% |
| best | 129 MB/s | 367 MB/s (2.8×) | 884 MB/s (6.8×) | 8.48% | 8.63% |
压缩等级通过WithEncoderLevel()选项指定,目前仅支持预定义档位。
3.5 未来兼容性承诺
- 压缩效率与速度未来可能变化;目标是将默认效率保持在原生 zstd level 3 附近;
- 不要对压缩输出做哈希相似度校验——编码不应被假定保持不变;
- 相同代码版本下 Encoder 输出可保证一致;未来可能出现破坏性模式,但不会在未显式开启选项的情况下启用;
- 该编码器不设计为(且大概率永远不)输出与参考编码器完全一致的比特流;
- 文档同时提醒:cgo 解压器(DataDog zstd)在无效输入上报错不完整、省略错误检查、忽略校验和,且似乎忽略拼接流(尽管拼接流属于规范的一部分),这也是选择纯 Go 实现的一个考量。
3.6 小块压缩:EncodeAll
EncodeAll(src, dst []byte) []byte将src全部编码并追加到dst后返回。该函数可被并发调用,每次调用只在调用方所在 goroutine 上执行。编码后的块可以拼接,结果等价于组合后的输入流;EncodeAll压缩的数据可用 Decoder 的流式或DecodeAll方式解压。
小块压缩应特别注意复用 encoder,暖机后可做到零分配;若提供容量足够的 dst 缓冲,可做到完全零分配:
import "github.com/klauspost/compress/zstd" // Create a writer that caches compressors. // For this operation type we supply a nil Reader. var encoder, _ = zstd.NewWriter(nil) // Compress a buffer. // If you have a destination buffer, the allocation in the call can also be eliminated. func Compress(src []byte) []byte { return encoder.EncodeAll(src, make([]byte, 0, len(src))) }用WithEncoderConcurrency(n)可控制最大并发编码数。Encoder 同时用于流式与独立块编码是安全的。
四、解压器(Decompressor)详解
4.1 稳定性与模糊测试
解压器状态同样为STABLE。该库持续进行模糊测试,核心目的是保证任何输入都无法使解码器崩溃或超出其限制运行。
4.2 流式解压
解压器通过Decoder访问,两种主要用法分别面向大数据流与小内存缓冲。流式解压的简单示例:
import "github.com/klauspost/compress/zstd" func Decompress(in io.Reader, out io.Writer) error { d, err := zstd.NewReader(in) if err != nil { return err } defer d.Close() // Copy content... _, err = io.Copy(out, d) return err }重要:不再需要 Reader 时必须调用Close(),以停止默认设置下运行的 goroutine。goroutine 会在返回错误(包括流末尾的io.EOF)后退出。
流默认以 4 个异步阶段并发解码以获得最佳吞吐;若希望同步解压(仅在请求数据时解压),使用WithDecoderConcurrency(1)。
4.3 缓冲解压:DecodeAll
import "github.com/klauspost/compress/zstd" // Create a reader that caches decompressors. // For this operation type we supply a nil Reader. var decoder, _ = zstd.NewReader(nil, zstd.WithDecoderConcurrency(0)) // Decompress a buffer. We don't supply a destination buffer, // so it will be allocated by the decoder. func Decompress(src []byte) ([]byte, error) { return decoder.DecodeAll(src, nil) }- Decoder 可并发解压多个缓冲,默认创建 4 个解压器;
- 用
WithDecoderConcurrency(n)调节并发操作数,WithDecoderConcurrency(0)表示创建 GOMAXPROCS 个解压器。
4.4 字典压缩(Dictionaries)
支持解压使用字典压缩的数据。字典通过zstd --train命令生成,包含解码器的初始状态:
- 添加字典:
WithDecoderDicts(dicts ...[]byte),可一次添加多个;数据中指定了对应字典时会自动使用;复用的 Decoder 仍保留已注册字典;多个同 ID 字典注册时以最后一个为准; - 压缩侧启用:
WithEncoderDict(dict []byte),仅使用一个字典,即使它不改善压缩比也可能被使用; - 解压时必须使用与压缩时相同的字典;字典应基于相似数据构建,否则输出可能比不用字典还略大;
- 当前使用字典压缩存在固定的启动性能开销,实现时务必自行基准测试。
4.5 零分配运行与资源释放
解码器设计为暖机后零分配运行,因此应长期持有并复用 decoder:
- 复用流式解码器:
Reset(r io.Reader) error切换到另一条流;即使前一条流解码失败,decoder 也可安全复用; - 释放资源必须调用
Close();调用后不可再复用,但所有运行中的 goroutine 都会停止; - 解码小缓冲时,可传入长度为 0、容量为期望值的目标切片,从而避免不必要的分配。
4.6 解码并发模型
- 缓冲解码在同一 goroutine 上同步执行,不做内部并发,但可并发解码多个缓冲(用
WithDecoderConcurrency(n)限制); - 流式解码器创建 4 类 goroutine:① 读取输入并切分为块;② 字面量解压;③ 序列解压;④ 输出流重建。这使解码器具备"预读"能力;
- 流的并发级别决定了解压提前多少块开始;由于块与前一块的输出强相关,流式解码的并发有限,实践中通常只能有效利用约 3 个核心。
4.7 解码基准
以下基准运行于 AMD Ryzen 9 3950X 16 核处理器,使用 AMD64 汇编(前两条为流式解码,其余为小输入并发解码):
BenchmarkDecoderSilesia-32 5 206878840 ns/op 1024.50 MB/s 49808 B/op 43 allocs/op BenchmarkDecoderEnwik9-32 1 1271809000 ns/op 786.28 MB/s 72048 B/op 52 allocs/op Concurrent blocks, performance: BenchmarkDecoder_DecodeAllParallel/kppkn.gtb.zst-32 67356 17857 ns/op 10321.96 MB/s 22.48 pct 102 B/op 0 allocs/op BenchmarkDecoder_DecodeAllParallel/geo.protodata.zst-32 266656 4421 ns/op 26823.21 MB/s 11.89 pct 19 B/op 0 allocs/op BenchmarkDecoder_DecodeAllParallel/plrabn12.txt.zst-32 20992 56842 ns/op 8477.17 MB/s 39.90 pct 754 B/op 0 allocs/op BenchmarkDecoder_DecodeAllParallel/lcet10.txt.zst-32 27456 43932 ns/op 9714.01 MB/s 33.27 pct 524 B/op 0 allocs/op BenchmarkDecoder_DecodeAllParallel/asyoulik.txt.zst-32 78432 15047 ns/op 8319.15 MB/s 40.34 pct 66 B/op 0 allocs/op BenchmarkDecoder_DecodeAllParallel/alice29.txt.zst-32 65800 18436 ns/op 8249.63 MB/s 37.75 pct 88 B/op 0 allocs/op BenchmarkDecoder_DecodeAllParallel/html_x_4.zst-32 102993 11523 ns/op 35546.09 MB/s 3.637 pct 143 B/op 0 allocs/op BenchmarkDecoder_DecodeAllParallel/paper-100k.pdf.zst-32 1000000 1070 ns/op 95720.98 MB/s 80.53 pct 3 B/op 0 allocs/op BenchmarkDecoder_DecodeAllParallel/fireworks.jpeg.zst-32 749802 1752 ns/op 70272.35 MB/s 100.0 pct 5 B/op 0 allocs/op BenchmarkDecoder_DecodeAllParallel/urls.10K.zst-32 22640 52934 ns/op 13263.37 MB/s 26.25 pct 1014 B/op 0 allocs/op BenchmarkDecoder_DecodeAllParallel/html.zst-32 226412 5232 ns/op 19572.27 MB/s 14.49 pct 20 B/op 0 allocs/op BenchmarkDecoder_DecodeAllParallel/comp-data.bin.zst-32 923041 1276 ns/op 3194.71 MB/s 31.26 pct 0 B/op 0 allocs/op文档注明上述数据反映约 2022 年 5 月的性能,可能已过时。
五、在 ZIP 文件中使用 zstd
可以在 zip 归档内对单个文件使用 zstandard 压缩,虽未被广泛支持,但对内部文件很有用。做法是注册对应的压缩器/解压器,文档强烈建议在单个 zip Reader/Writer 上注册(de)compressors,而不要使用全局注册函数——来自不同包的两处全局注册会引发 panic。较好的实践是只保留一个压缩器与一个解压器实例,它们可被多个 zip 文件并发复用,单实例还能复用部分内部资源。
六、Tempo 仓库中的真实调用链(源码佐证)
github.com/klauspost/compress/zstd并非仅仅被 vendored 进来,而是 Tempo 多个关键路径的运行时依赖。以下调用点可在仓库源码中直接验证。
6.1 块后端编解码器(tempodb/backend/compression.go)
Tempo 的块后端抽象定义了Codec接口,ZstdCodec是其中一个实现,直接使用本库:
- 编解码器通过
sync.Pool缓存*zstd.Encoder与*zstd.Decoder实例,实现实例复用; Encode使用zstd.NewWriter(nil, zstd.WithEncoderConcurrency(1))创建流式 writer 并调用EncodeAll(src, dst)——这正是文档中"小块压缩优先使用 EncodeAll、用WithEncoderConcurrency(1)关闭流式异步"两个建议的组合;Decode使用zstd.NewReader(nil, zstd.WithDecoderConcurrency(0))并调用DecodeAll(buf, nil),其中WithDecoderConcurrency(0)表示创建 GOMAXPROCS 个解码器,与文档 4.3 节所述一致。
// 摘自 tempodb/backend/compression.go func (c *ZstdCodec) Encode(src, dst []byte) ([]byte, error) { e, _ := c.encoders.Get().(*zstd.Encoder) if e == nil { var err error e, err = zstd.NewWriter(nil, zstd.WithEncoderConcurrency(1)) if err != nil { return nil, err } } defer c.encoders.Put(e) return e.EncodeAll(src, dst), nil } func (c *ZstdCodec) Decode(buf []byte) ([]byte, error) { d, _ := c.decoders.Get().(*zstd.Decoder) if d == nil { var err error d, err = zstd.NewReader(nil, zstd.WithDecoderConcurrency(0)) if err != nil { return nil, err } } defer c.decoders.Put(d) return d.DecodeAll(buf, nil) }6.2 租户索引(tempodb/backend/tenantindex.go)
tenantindex.go中定义了全局Zstd = &ZstdCodec{},其 protobuf 序列化路径marshalPb/unmarshalPb直接调用Zstd.Encode/Zstd.Decode对租户索引做 zstd 压缩,解压失败时返回error decoding zstd包装错误。这是"Decoder 必须与 Encoder 严格配对使用"在生产环境中的典型体现。
6.3 Kafka 生产者压缩(pkg/ingest/config.go)
Tempo 的 ingest 配置允许为 Kafka 生产者选择压缩编解码器,zstd是受支持取值之一:
- 配置项
producer-compression,合法值包括none, gzip, snappy, lz4, zstd; - 命中
compressionZstd分支时调用kgo.ZstdCompression()返回 Kafka 客户端压缩选项,底层即依赖本仓库的 zstd 实现; - 对应错误类型
ErrInvalidProducerCompression会在配置了非法编解码器名时抛出。
6.4 Parquet 列压缩(tempodb/encoding/vparquet5/schema.go)
在 vparquet5 的列式存储 schema 构建中,Tempo 对指定列使用parquet:",zstd,optional"的结构体标签,将列压缩方式改为 zstd(同时移除字典编码),进一步印证 zstd 在 Tempo 存储链路中的广泛使用。
七、工程实践要点总结
结合官方文档与 Tempo 源码,可沉淀出以下可直接落地的经验:
- 小块用
EncodeAll、大流用NewWriter,两者通过nilwriter 创建的实例可复用且支持并发; - 长期持有并复用 encoder/decoder(Tempo 用
sync.Pool实现),配合预分配 dst 缓冲可达到暖机后零分配; - 显式设置并发:流式编码默认 2 goroutine 的"轻并发"可能在版本演进中变化,需要确定性并发时显式传
WithEncoderConcurrency(n);流式解码默认 4 异步阶段、实际约利用 3 核; - 必须成对使用字典:同 ID 多字典取最后注册者;字典应基于相似数据训练,并留意固定启动开销;
- 资源释放纪律:编码失败也要
Close();decoderClose()后不可复用;复用流目标用Reset(); - 不要依赖比特级兼容:编码输出与参考实现不保证一致,禁用压缩输出哈希做相似性校验。
上述所有调用点均可在仓库源码中逐一核对,例如 tempodb/backend/compression.go、tempodb/backend/tenantindex.go、pkg/ingest/config.go 与 tempodb/encoding/vparquet5/schema.go。
【免费下载链接】tempoGrafana Tempo is a high volume, minimal dependency distributed tracing backend.项目地址: https://gitcode.com/GitHub_Trending/tempo1/tempo
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考