Skip to content

Commit ce45805

Browse files
felixgensrip-dd
authored andcommitted
fix(profiler): support gzip compression mode [#4696 backport]
Backport #4696 to the v2.7.x release branch Adds a functional gzip compression mode for the profiler via `DD_PROFILING_DEBUG_COMPRESSION_SETTINGS=gzip`. This was already kinda working, but not really, hence the "fix" title. - Defaults `gzip` to gzip-6, matching the gzip level legacy mode applies when it compresses profiler-produced profiles. - Adds gzip-to-gzip recompression so gzip mode works for runtime profiles that are already gzip-compressed at gzip-1. - Adds tests for gzip mode and gzip recompression. This is done for incident-53386, where zstd caused a memory usage increase. The gzip mode provides an alternative to the default zstd compression while keeping profiler uploads compressed.
1 parent dc57de6 commit ce45805

2 files changed

Lines changed: 67 additions & 2 deletions

File tree

profiler/compression.go

Lines changed: 56 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -48,11 +48,20 @@ func compressionStrategy(pt ProfileType, isDelta bool, config string) (compressi
4848
if config == "" || config == "legacy" {
4949
return legacyCompressionStrategy(pt, isDelta)
5050
}
51-
algorithm, levelStr, _ := strings.Cut(config, "-")
51+
algorithmStr, levelStr, _ := strings.Cut(config, "-")
52+
algorithm := compressionAlgorithm(algorithmStr)
5253
// Don't bother checking the error. We'll get zero which represents the
5354
// default, and we we assume this is only going to get used internally
5455
level, _ := strconv.Atoi(levelStr)
55-
return inputCompression(pt, isDelta), compression{algorithm: compressionAlgorithm(algorithm), level: level}
56+
if levelStr == "" {
57+
switch algorithm {
58+
case compressionAlgorithmGzip:
59+
level = gzip6Compression.level
60+
case compressionAlgorithmZstd:
61+
level = zstdCompression.level
62+
}
63+
}
64+
return inputCompression(pt, isDelta), compression{algorithm: algorithm, level: level}
5665
}
5766

5867
// inputCompression maps the given profile type and isDelta flavor to the
@@ -172,6 +181,14 @@ func (b *compressionPipelineBuilder) Build(in compression, out compression) (com
172181
return b.getZstdEncoder(getZstdLevelOrDefault(out.level))
173182
}
174183

184+
if in.algorithm == compressionAlgorithmGzip && out.algorithm == compressionAlgorithmGzip {
185+
gzipOut, err := kgzip.NewWriterLevel(nil, out.level)
186+
if err != nil {
187+
return nil, err
188+
}
189+
return newGzipRecompressor(gzipOut), nil
190+
}
191+
175192
if in.algorithm == compressionAlgorithmGzip && out.algorithm == compressionAlgorithmZstd {
176193
encoder, err := b.getZstdEncoder(getZstdLevelOrDefault(out.level))
177194
if err != nil {
@@ -213,6 +230,43 @@ func (r *passthroughCompressor) Close() error {
213230
return nil
214231
}
215232

233+
func newGzipRecompressor(gzipOut *kgzip.Writer) *gzipRecompressor {
234+
return &gzipRecompressor{gzipOut: gzipOut, err: make(chan error)}
235+
}
236+
237+
type gzipRecompressor struct {
238+
// err synchronizes finishing writes after closing pw and reports any
239+
// error during recompression
240+
err chan error
241+
pw io.WriteCloser
242+
gzipOut *kgzip.Writer
243+
}
244+
245+
func (r *gzipRecompressor) Reset(w io.Writer) {
246+
r.gzipOut.Reset(w)
247+
pr, pw := io.Pipe()
248+
go func() {
249+
gzr, err := kgzip.NewReader(pr)
250+
if err != nil {
251+
r.err <- err
252+
return
253+
}
254+
_, err = io.Copy(r.gzipOut, gzr)
255+
r.err <- err
256+
}()
257+
r.pw = pw
258+
}
259+
260+
func (r *gzipRecompressor) Write(p []byte) (int, error) {
261+
return r.pw.Write(p)
262+
}
263+
264+
func (r *gzipRecompressor) Close() error {
265+
r.pw.Close()
266+
err := <-r.err
267+
return cmp.Or(err, r.gzipOut.Close())
268+
}
269+
216270
func newZstdRecompressor(encoder *sharedZstdEncoder) *zstdRecompressor {
217271
return &zstdRecompressor{zstdOut: encoder, err: make(chan error)}
218272
}

profiler/compression_test.go

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,8 @@ func TestNewCompressionPipeline(t *testing.T) {
3737
{noCompression, noCompression, plainData, plainData},
3838
{gzip1Compression, gzip1Compression, gzip1Data, gzip1Data},
3939
{gzip6Compression, gzip6Compression, gzip6Data, gzip6Data},
40+
{gzip1Compression, gzip6Compression, gzip1Data, gzip6Data},
41+
{gzip6Compression, gzip1Compression, gzip6Data, gzip1Data},
4042
{gzip1Compression, zstdCompression, gzip1Data, zstdData},
4143
{gzip6Compression, zstdCompression, gzip6Data, zstdData},
4244
}
@@ -99,6 +101,15 @@ func TestDebugCompressionEnv(t *testing.T) {
99101
require.NoError(t, err)
100102
})
101103

104+
t.Run("explicit-gzip-already-gzipped-input", func(t *testing.T) {
105+
t.Setenv("DD_PROFILING_DEBUG_COMPRESSION_SETTINGS", "gzip")
106+
p := startTestProfiler(t, 1, WithProfileTypes(CPUProfile), WithPeriod(time.Millisecond)).ReceiveProfile(t)
107+
r, err := gzip.NewReader(bytes.NewReader(p.attachments["cpu.pprof"]))
108+
require.NoError(t, err)
109+
_, err = io.Copy(io.Discard, r)
110+
require.NoError(t, err)
111+
})
112+
102113
t.Run("zstd-delta", func(t *testing.T) {
103114
t.Setenv("DD_PROFILING_DEBUG_COMPRESSION_SETTINGS", "zstd-3")
104115
p := startTestProfiler(t, 1, WithProfileTypes(CPUProfile, HeapProfile), WithPeriod(time.Millisecond)).ReceiveProfile(t)

0 commit comments

Comments
 (0)