Coverage Report

Created: 2026-08-04 13:44

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/util/block_compression.cpp
Line
Count
Source
1
// Licensed to the Apache Software Foundation (ASF) under one
2
// or more contributor license agreements.  See the NOTICE file
3
// distributed with this work for additional information
4
// regarding copyright ownership.  The ASF licenses this file
5
// to you under the Apache License, Version 2.0 (the
6
// "License"); you may not use this file except in compliance
7
// with the License.  You may obtain a copy of the License at
8
//
9
//   http://www.apache.org/licenses/LICENSE-2.0
10
//
11
// Unless required by applicable law or agreed to in writing,
12
// software distributed under the License is distributed on an
13
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14
// KIND, either express or implied.  See the License for the
15
// specific language governing permissions and limitations
16
// under the License.
17
18
#include "util/block_compression.h"
19
20
#include <bzlib.h>
21
#include <gen_cpp/parquet_types.h>
22
#include <gen_cpp/segment_v2.pb.h>
23
#include <glog/logging.h>
24
25
#include <exception>
26
// Only used on x86 or x86_64
27
#if defined(__x86_64__) || defined(_M_X64) || defined(i386) || defined(__i386__) || \
28
        defined(__i386) || defined(_M_IX86)
29
#include <libdeflate.h>
30
#endif
31
#include <brotli/decode.h>
32
#include <glog/log_severity.h>
33
#include <glog/logging.h>
34
#include <lz4/lz4.h>
35
#include <lz4/lz4frame.h>
36
#include <lz4/lz4hc.h>
37
#include <snappy/snappy-sinksource.h>
38
#include <snappy/snappy.h>
39
#include <zconf.h>
40
#include <zlib.h>
41
#include <zstd.h>
42
#include <zstd_errors.h>
43
44
#include <algorithm>
45
#include <cstdint>
46
#include <limits>
47
#include <mutex>
48
#include <orc/Exceptions.hh>
49
#include <ostream>
50
#include <unordered_map>
51
52
#include "absl/strings/substitute.h"
53
#include "common/config.h"
54
#include "common/factory_creator.h"
55
#include "exec/common/endian.h"
56
#include "runtime/thread_context.h"
57
#include "util/decompressor.h"
58
#include "util/defer_op.h"
59
#include "util/faststring.h"
60
61
namespace orc {
62
/**
63
 * Decompress the bytes in to the output buffer.
64
 * @param inputAddress the start of the input
65
 * @param inputLimit one past the last byte of the input
66
 * @param outputAddress the start of the output buffer
67
 * @param outputLimit one past the last byte of the output buffer
68
 * @result the number of bytes decompressed
69
 */
70
uint64_t lzoDecompress(const char* inputAddress, const char* inputLimit, char* outputAddress,
71
                       char* outputLimit);
72
} // namespace orc
73
74
namespace doris {
75
76
// exception safe
77
Status BlockCompressionCodec::compress(const std::vector<Slice>& inputs, size_t uncompressed_size,
78
569
                                       faststring* output) {
79
569
    faststring buf;
80
    // we compute total size to avoid more memory copy
81
569
    buf.reserve(uncompressed_size);
82
634
    for (auto& input : inputs) {
83
634
        buf.append(input.data, input.size);
84
634
    }
85
569
    return compress(buf, output);
86
569
}
87
88
392k
bool BlockCompressionCodec::exceed_max_compress_len(size_t uncompressed_size) {
89
392k
    return uncompressed_size > std::numeric_limits<int32_t>::max();
90
392k
}
91
92
class Lz4BlockCompression : public BlockCompressionCodec {
93
private:
94
    class Context {
95
        ENABLE_FACTORY_CREATOR(Context);
96
97
    public:
98
30
        Context() : ctx(nullptr) {
99
30
            SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
100
30
                    ExecEnv::GetInstance()->block_compression_mem_tracker());
101
30
            buffer = std::make_unique<faststring>();
102
30
        }
103
        LZ4_stream_t* ctx;
104
        std::unique_ptr<faststring> buffer;
105
8
        ~Context() {
106
8
            SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
107
8
                    ExecEnv::GetInstance()->block_compression_mem_tracker());
108
8
            if (ctx) {
109
8
                LZ4_freeStream(ctx);
110
8
            }
111
8
            buffer.reset();
112
8
        }
113
    };
114
115
public:
116
98.3k
    static Lz4BlockCompression* instance() {
117
98.3k
        static Lz4BlockCompression s_instance;
118
98.3k
        return &s_instance;
119
98.3k
    }
120
5
    ~Lz4BlockCompression() override {
121
5
        SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
122
5
                ExecEnv::GetInstance()->block_compression_mem_tracker());
123
5
        _ctx_pool.clear();
124
5
    }
125
126
48.5k
    Status compress(const Slice& input, faststring* output) override {
127
48.5k
        if (input.size > LZ4_MAX_INPUT_SIZE) {
128
0
            return Status::InvalidArgument(
129
0
                    "LZ4 not support those case(input.size>LZ4_MAX_INPUT_SIZE), maybe you should "
130
0
                    "change "
131
0
                    "fragment_transmission_compression_codec to snappy, input.size={}, "
132
0
                    "LZ4_MAX_INPUT_SIZE={}",
133
0
                    input.size, LZ4_MAX_INPUT_SIZE);
134
0
        }
135
136
48.5k
        std::unique_ptr<Context> context;
137
48.5k
        RETURN_IF_ERROR(_acquire_compression_ctx(context));
138
48.5k
        bool compress_failed = false;
139
48.5k
        Defer defer {[&] {
140
48.5k
            if (!compress_failed) {
141
48.5k
                _release_compression_ctx(std::move(context));
142
48.5k
            }
143
48.5k
        }};
144
145
48.5k
        try {
146
48.5k
            Slice compressed_buf;
147
48.5k
            size_t max_len = max_compressed_len(input.size);
148
48.5k
            if (max_len > MAX_COMPRESSION_BUFFER_SIZE_FOR_REUSE) {
149
                // use output directly
150
4
                output->resize(max_len);
151
4
                compressed_buf.data = reinterpret_cast<char*>(output->data());
152
4
                compressed_buf.size = max_len;
153
48.5k
            } else {
154
                // reuse context buffer if max_len <= MAX_COMPRESSION_BUFFER_FOR_REUSE
155
48.5k
                {
156
                    // context->buffer is resuable between queries, should accouting to
157
                    // global tracker.
158
48.5k
                    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
159
48.5k
                            ExecEnv::GetInstance()->block_compression_mem_tracker());
160
48.5k
                    context->buffer->resize(max_len);
161
48.5k
                }
162
48.5k
                compressed_buf.data = reinterpret_cast<char*>(context->buffer->data());
163
48.5k
                compressed_buf.size = max_len;
164
48.5k
            }
165
166
            // input.size is aready checked before;
167
            // compressed_buf.size is got from max_compressed_len, which is
168
            // the return value of LZ4_compressBound, so it is safe to cast to int
169
48.5k
            size_t compressed_len = LZ4_compress_fast_continue(
170
48.5k
                    context->ctx, input.data, compressed_buf.data, static_cast<int>(input.size),
171
48.5k
                    static_cast<int>(compressed_buf.size), ACCELARATION);
172
48.5k
            if (compressed_len == 0) {
173
0
                compress_failed = true;
174
0
                return Status::InvalidArgument("Output buffer's capacity is not enough, size={}",
175
0
                                               compressed_buf.size);
176
0
            }
177
48.5k
            output->resize(compressed_len);
178
48.5k
            if (max_len <= MAX_COMPRESSION_BUFFER_SIZE_FOR_REUSE) {
179
48.5k
                output->assign_copy(reinterpret_cast<uint8_t*>(compressed_buf.data),
180
48.5k
                                    compressed_len);
181
48.5k
            }
182
48.5k
        } catch (...) {
183
            // Do not set compress_failed to release context
184
0
            DCHECK(!compress_failed);
185
0
            return Status::InternalError("Fail to do LZ4Block compress due to exception");
186
0
        }
187
48.5k
        return Status::OK();
188
48.5k
    }
189
190
51.6k
    Status decompress(const Slice& input, Slice* output) override {
191
51.6k
        auto decompressed_len = LZ4_decompress_safe(
192
51.6k
                input.data, output->data, cast_set<int>(input.size), cast_set<int>(output->size));
193
51.6k
        if (decompressed_len < 0) {
194
5
            return Status::InternalError("fail to do LZ4 decompress, error={}", decompressed_len);
195
5
        }
196
51.6k
        output->size = decompressed_len;
197
51.6k
        return Status::OK();
198
51.6k
    }
199
200
48.5k
    size_t max_compressed_len(size_t len) override { return LZ4_compressBound(cast_set<int>(len)); }
201
202
private:
203
    // reuse LZ4 compress stream
204
48.5k
    Status _acquire_compression_ctx(std::unique_ptr<Context>& out) {
205
48.5k
        std::lock_guard<std::mutex> l(_ctx_mutex);
206
48.5k
        if (_ctx_pool.empty()) {
207
30
            std::unique_ptr<Context> localCtx = Context::create_unique();
208
30
            if (localCtx.get() == nullptr) {
209
0
                return Status::InvalidArgument("new LZ4 context error");
210
0
            }
211
30
            localCtx->ctx = LZ4_createStream();
212
30
            if (localCtx->ctx == nullptr) {
213
0
                return Status::InvalidArgument("LZ4_createStream error");
214
0
            }
215
30
            out = std::move(localCtx);
216
30
            return Status::OK();
217
30
        }
218
48.5k
        out = std::move(_ctx_pool.back());
219
48.5k
        _ctx_pool.pop_back();
220
48.5k
        return Status::OK();
221
48.5k
    }
222
48.5k
    void _release_compression_ctx(std::unique_ptr<Context> context) {
223
48.5k
        DCHECK(context);
224
48.5k
        LZ4_resetStream(context->ctx);
225
48.5k
        std::lock_guard<std::mutex> l(_ctx_mutex);
226
48.5k
        _ctx_pool.push_back(std::move(context));
227
48.5k
    }
228
229
private:
230
    mutable std::mutex _ctx_mutex;
231
    mutable std::vector<std::unique_ptr<Context>> _ctx_pool;
232
    static const int32_t ACCELARATION = 1;
233
};
234
235
class HadoopLz4BlockCompression : public Lz4BlockCompression {
236
public:
237
3
    HadoopLz4BlockCompression() {
238
3
        Status st = Decompressor::create_decompressor(CompressType::LZ4BLOCK, &_decompressor);
239
3
        if (!st.ok()) {
240
0
            throw Exception(Status::FatalError(
241
0
                    "HadoopLz4BlockCompression construction failed. status = {}", st));
242
0
        }
243
3
    }
244
245
1
    ~HadoopLz4BlockCompression() override = default;
246
247
2.25k
    static HadoopLz4BlockCompression* instance() {
248
2.25k
        static HadoopLz4BlockCompression s_instance;
249
2.25k
        return &s_instance;
250
2.25k
    }
251
252
    // hadoop use block compression for lz4
253
    // https://github.com/apache/hadoop/blob/trunk/hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-nativetask/src/main/native/src/codec/Lz4Codec.cc
254
277
    Status compress(const Slice& input, faststring* output) override {
255
        // be same with hadop https://github.com/apache/hadoop/blob/trunk/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/Lz4Codec.java
256
277
        size_t lz4_block_size = config::lz4_compression_block_size;
257
277
        size_t overhead = lz4_block_size / 255 + 16;
258
277
        size_t max_input_size = lz4_block_size - overhead;
259
260
277
        size_t data_len = input.size;
261
277
        char* data = input.data;
262
277
        std::vector<OwnedSlice> buffers;
263
277
        size_t out_len = 0;
264
265
600
        while (data_len > 0) {
266
323
            size_t input_size = std::min(data_len, max_input_size);
267
323
            Slice input_slice(data, input_size);
268
323
            faststring output_data;
269
323
            RETURN_IF_ERROR(Lz4BlockCompression::compress(input_slice, &output_data));
270
323
            out_len += output_data.size();
271
323
            buffers.push_back(output_data.build());
272
323
            data += input_size;
273
323
            data_len -= input_size;
274
323
        }
275
276
        // hadoop block compression: umcompressed_length | compressed_length1 | compressed_data1 | compressed_length2 | compressed_data2 | ...
277
277
        size_t total_output_len = 4 + 4 * buffers.size() + out_len;
278
277
        output->resize(total_output_len);
279
277
        char* output_buffer = (char*)output->data();
280
277
        BigEndian::Store32(output_buffer, cast_set<uint32_t>(input.get_size()));
281
277
        output_buffer += 4;
282
323
        for (const auto& buffer : buffers) {
283
323
            auto slice = buffer.slice();
284
323
            BigEndian::Store32(output_buffer, cast_set<uint32_t>(slice.get_size()));
285
323
            output_buffer += 4;
286
323
            memcpy(output_buffer, slice.get_data(), slice.get_size());
287
323
            output_buffer += slice.get_size();
288
323
        }
289
290
277
        DCHECK_EQ(output_buffer - (char*)output->data(), total_output_len);
291
292
277
        return Status::OK();
293
277
    }
294
295
2.05k
    Status decompress(const Slice& input, Slice* output) override {
296
2.05k
        size_t input_bytes_read = 0;
297
2.05k
        size_t decompressed_len = 0;
298
2.05k
        size_t more_input_bytes = 0;
299
2.05k
        size_t more_output_bytes = 0;
300
2.05k
        bool stream_end = false;
301
2.05k
        auto st = _decompressor->decompress((uint8_t*)input.data, cast_set<uint32_t>(input.size),
302
2.05k
                                            &input_bytes_read, (uint8_t*)output->data,
303
2.05k
                                            cast_set<uint32_t>(output->size), &decompressed_len,
304
2.05k
                                            &stream_end, &more_input_bytes, &more_output_bytes);
305
        //try decompress use hadoopLz4 ,if failed fall back lz4.
306
2.05k
        return (st != Status::OK() || stream_end != true)
307
2.05k
                       ? Lz4BlockCompression::decompress(input, output)
308
2.05k
                       : Status::OK();
309
2.05k
    }
310
311
private:
312
    std::unique_ptr<Decompressor> _decompressor;
313
};
314
// Used for LZ4 frame format, decompress speed is two times faster than LZ4.
315
class Lz4fBlockCompression : public BlockCompressionCodec {
316
private:
317
    class CContext {
318
        ENABLE_FACTORY_CREATOR(CContext);
319
320
    public:
321
1
        CContext() : ctx(nullptr) {
322
1
            SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
323
1
                    ExecEnv::GetInstance()->block_compression_mem_tracker());
324
1
            buffer = std::make_unique<faststring>();
325
1
        }
326
        LZ4F_compressionContext_t ctx;
327
        std::unique_ptr<faststring> buffer;
328
1
        ~CContext() {
329
1
            SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
330
1
                    ExecEnv::GetInstance()->block_compression_mem_tracker());
331
1
            if (ctx) {
332
1
                LZ4F_freeCompressionContext(ctx);
333
1
            }
334
1
            buffer.reset();
335
1
        }
336
    };
337
    class DContext {
338
        ENABLE_FACTORY_CREATOR(DContext);
339
340
    public:
341
6
        DContext() : ctx(nullptr) {}
342
        LZ4F_decompressionContext_t ctx;
343
6
        ~DContext() {
344
6
            if (ctx) {
345
6
                LZ4F_freeDecompressionContext(ctx);
346
6
            }
347
6
        }
348
    };
349
350
public:
351
26.2k
    static Lz4fBlockCompression* instance() {
352
26.2k
        static Lz4fBlockCompression s_instance;
353
26.2k
        return &s_instance;
354
26.2k
    }
355
1
    ~Lz4fBlockCompression() {
356
1
        SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
357
1
                ExecEnv::GetInstance()->block_compression_mem_tracker());
358
1
        _ctx_c_pool.clear();
359
1
        _ctx_d_pool.clear();
360
1
    }
361
362
5
    Status compress(const Slice& input, faststring* output) override {
363
5
        std::vector<Slice> inputs {input};
364
5
        return compress(inputs, input.size, output);
365
5
    }
366
367
    Status compress(const std::vector<Slice>& inputs, size_t uncompressed_size,
368
23.6k
                    faststring* output) override {
369
23.6k
        return _compress(inputs, uncompressed_size, output);
370
23.6k
    }
371
372
8.38k
    Status decompress(const Slice& input, Slice* output) override {
373
8.38k
        return _decompress(input, output);
374
8.38k
    }
375
376
23.6k
    size_t max_compressed_len(size_t len) override {
377
23.6k
        return std::max(LZ4F_compressBound(len, &_s_preferences),
378
23.6k
                        LZ4F_compressFrameBound(len, &_s_preferences));
379
23.6k
    }
380
381
private:
382
    Status _compress(const std::vector<Slice>& inputs, size_t uncompressed_size,
383
23.6k
                     faststring* output) {
384
23.6k
        std::unique_ptr<CContext> context;
385
23.6k
        RETURN_IF_ERROR(_acquire_compression_ctx(context));
386
23.6k
        bool compress_failed = false;
387
23.6k
        Defer defer {[&] {
388
23.6k
            if (!compress_failed) {
389
23.6k
                _release_compression_ctx(std::move(context));
390
23.6k
            }
391
23.6k
        }};
392
393
23.6k
        try {
394
23.6k
            Slice compressed_buf;
395
23.6k
            size_t max_len = max_compressed_len(uncompressed_size);
396
23.6k
            if (max_len > MAX_COMPRESSION_BUFFER_SIZE_FOR_REUSE) {
397
                // use output directly
398
0
                output->resize(max_len);
399
0
                compressed_buf.data = reinterpret_cast<char*>(output->data());
400
0
                compressed_buf.size = max_len;
401
23.6k
            } else {
402
23.6k
                {
403
23.6k
                    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
404
23.6k
                            ExecEnv::GetInstance()->block_compression_mem_tracker());
405
                    // reuse context buffer if max_len <= MAX_COMPRESSION_BUFFER_FOR_REUSE
406
23.6k
                    context->buffer->resize(max_len);
407
23.6k
                }
408
23.6k
                compressed_buf.data = reinterpret_cast<char*>(context->buffer->data());
409
23.6k
                compressed_buf.size = max_len;
410
23.6k
            }
411
412
23.6k
            auto wbytes = LZ4F_compressBegin(context->ctx, compressed_buf.data, compressed_buf.size,
413
23.6k
                                             &_s_preferences);
414
23.6k
            if (LZ4F_isError(wbytes)) {
415
0
                compress_failed = true;
416
0
                return Status::InvalidArgument("Fail to do LZ4F compress begin, res={}",
417
0
                                               LZ4F_getErrorName(wbytes));
418
0
            }
419
23.6k
            size_t offset = wbytes;
420
24.2k
            for (auto input : inputs) {
421
24.2k
                wbytes = LZ4F_compressUpdate(context->ctx, compressed_buf.data + offset,
422
24.2k
                                             compressed_buf.size - offset, input.data, input.size,
423
24.2k
                                             nullptr);
424
24.2k
                if (LZ4F_isError(wbytes)) {
425
0
                    compress_failed = true;
426
0
                    return Status::InvalidArgument("Fail to do LZ4F compress update, res={}",
427
0
                                                   LZ4F_getErrorName(wbytes));
428
0
                }
429
24.2k
                offset += wbytes;
430
24.2k
            }
431
23.6k
            wbytes = LZ4F_compressEnd(context->ctx, compressed_buf.data + offset,
432
23.6k
                                      compressed_buf.size - offset, nullptr);
433
23.6k
            if (LZ4F_isError(wbytes)) {
434
0
                compress_failed = true;
435
0
                return Status::InvalidArgument("Fail to do LZ4F compress end, res={}",
436
0
                                               LZ4F_getErrorName(wbytes));
437
0
            }
438
23.6k
            offset += wbytes;
439
23.6k
            output->resize(offset);
440
23.6k
            if (max_len <= MAX_COMPRESSION_BUFFER_SIZE_FOR_REUSE) {
441
23.6k
                output->assign_copy(reinterpret_cast<uint8_t*>(compressed_buf.data), offset);
442
23.6k
            }
443
23.6k
        } catch (...) {
444
            // Do not set compress_failed to release context
445
0
            DCHECK(!compress_failed);
446
0
            return Status::InternalError("Fail to do LZ4F compress due to exception");
447
0
        }
448
449
23.6k
        return Status::OK();
450
23.6k
    }
451
452
8.38k
    Status _decompress(const Slice& input, Slice* output) {
453
8.38k
        bool decompress_failed = false;
454
8.38k
        std::unique_ptr<DContext> context;
455
8.38k
        RETURN_IF_ERROR(_acquire_decompression_ctx(context));
456
8.38k
        Defer defer {[&] {
457
8.38k
            if (!decompress_failed) {
458
8.37k
                _release_decompression_ctx(std::move(context));
459
8.37k
            }
460
8.38k
        }};
461
8.38k
        size_t input_size = input.size;
462
8.38k
        auto lres = LZ4F_decompress(context->ctx, output->data, &output->size, input.data,
463
8.38k
                                    &input_size, nullptr);
464
8.38k
        if (LZ4F_isError(lres)) {
465
0
            decompress_failed = true;
466
0
            return Status::InternalError("Fail to do LZ4F decompress, res={}",
467
0
                                         LZ4F_getErrorName(lres));
468
8.38k
        } else if (input_size != input.size) {
469
0
            decompress_failed = true;
470
0
            return Status::InvalidArgument(
471
0
                    absl::Substitute("Fail to do LZ4F decompress: trailing data left in "
472
0
                                     "compressed data, read=$0 vs given=$1",
473
0
                                     input_size, input.size));
474
8.38k
        } else if (lres != 0) {
475
5
            decompress_failed = true;
476
5
            return Status::InvalidArgument(
477
5
                    "Fail to do LZ4F decompress: expect more compressed data, expect={}", lres);
478
5
        }
479
8.37k
        return Status::OK();
480
8.38k
    }
481
482
private:
483
    // acquire a compression ctx from pool, release while finish compress,
484
    // delete if compression failed
485
23.6k
    Status _acquire_compression_ctx(std::unique_ptr<CContext>& out) {
486
23.6k
        std::lock_guard<std::mutex> l(_ctx_c_mutex);
487
23.6k
        if (_ctx_c_pool.empty()) {
488
1
            std::unique_ptr<CContext> localCtx = CContext::create_unique();
489
1
            if (localCtx.get() == nullptr) {
490
0
                return Status::InvalidArgument("failed to new LZ4F CContext");
491
0
            }
492
1
            auto res = LZ4F_createCompressionContext(&localCtx->ctx, LZ4F_VERSION);
493
1
            if (LZ4F_isError(res) != 0) {
494
0
                return Status::InvalidArgument(absl::Substitute(
495
0
                        "LZ4F_createCompressionContext error, res=$0", LZ4F_getErrorName(res)));
496
0
            }
497
1
            out = std::move(localCtx);
498
1
            return Status::OK();
499
1
        }
500
23.6k
        out = std::move(_ctx_c_pool.back());
501
23.6k
        _ctx_c_pool.pop_back();
502
23.6k
        return Status::OK();
503
23.6k
    }
504
23.6k
    void _release_compression_ctx(std::unique_ptr<CContext> context) {
505
23.6k
        DCHECK(context);
506
23.6k
        std::lock_guard<std::mutex> l(_ctx_c_mutex);
507
23.6k
        _ctx_c_pool.push_back(std::move(context));
508
23.6k
    }
509
510
8.38k
    Status _acquire_decompression_ctx(std::unique_ptr<DContext>& out) {
511
8.38k
        std::lock_guard<std::mutex> l(_ctx_d_mutex);
512
8.38k
        if (_ctx_d_pool.empty()) {
513
6
            std::unique_ptr<DContext> localCtx = DContext::create_unique();
514
6
            if (localCtx.get() == nullptr) {
515
0
                return Status::InvalidArgument("failed to new LZ4F DContext");
516
0
            }
517
6
            auto res = LZ4F_createDecompressionContext(&localCtx->ctx, LZ4F_VERSION);
518
6
            if (LZ4F_isError(res) != 0) {
519
0
                return Status::InvalidArgument(absl::Substitute(
520
0
                        "LZ4F_createDeompressionContext error, res=$0", LZ4F_getErrorName(res)));
521
0
            }
522
6
            out = std::move(localCtx);
523
6
            return Status::OK();
524
6
        }
525
8.37k
        out = std::move(_ctx_d_pool.back());
526
8.37k
        _ctx_d_pool.pop_back();
527
8.37k
        return Status::OK();
528
8.38k
    }
529
8.37k
    void _release_decompression_ctx(std::unique_ptr<DContext> context) {
530
8.37k
        DCHECK(context);
531
        // reset decompression context to avoid ERROR_maxBlockSize_invalid
532
8.37k
        LZ4F_resetDecompressionContext(context->ctx);
533
8.37k
        std::lock_guard<std::mutex> l(_ctx_d_mutex);
534
8.37k
        _ctx_d_pool.push_back(std::move(context));
535
8.37k
    }
536
537
private:
538
    static LZ4F_preferences_t _s_preferences;
539
540
    std::mutex _ctx_c_mutex;
541
    // LZ4F_compressionContext_t is a pointer so no copy here
542
    std::vector<std::unique_ptr<CContext>> _ctx_c_pool;
543
544
    std::mutex _ctx_d_mutex;
545
    std::vector<std::unique_ptr<DContext>> _ctx_d_pool;
546
};
547
548
LZ4F_preferences_t Lz4fBlockCompression::_s_preferences = {
549
        {LZ4F_max256KB, LZ4F_blockLinked, LZ4F_noContentChecksum, LZ4F_frame, 0ULL, 0U,
550
         LZ4F_noBlockChecksum},
551
        0,
552
        0u,
553
        0u,
554
        {0u, 0u, 0u}};
555
556
class Lz4HCBlockCompression : public BlockCompressionCodec {
557
private:
558
    class Context {
559
        ENABLE_FACTORY_CREATOR(Context);
560
561
    public:
562
4
        Context() : ctx(nullptr) {
563
4
            SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
564
4
                    ExecEnv::GetInstance()->block_compression_mem_tracker());
565
4
            buffer = std::make_unique<faststring>();
566
4
        }
567
        LZ4_streamHC_t* ctx;
568
        std::unique_ptr<faststring> buffer;
569
4
        ~Context() {
570
4
            SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
571
4
                    ExecEnv::GetInstance()->block_compression_mem_tracker());
572
4
            if (ctx) {
573
4
                LZ4_freeStreamHC(ctx);
574
4
            }
575
4
            buffer.reset();
576
4
        }
577
    };
578
579
public:
580
3
    static Lz4HCBlockCompression* instance() {
581
3
        static Lz4HCBlockCompression s_instance;
582
3
        return &s_instance;
583
3
    }
584
1
    Lz4HCBlockCompression() = default;
585
4
    explicit Lz4HCBlockCompression(int level) : _compression_level(level) {}
586
5
    ~Lz4HCBlockCompression() {
587
5
        SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
588
5
                ExecEnv::GetInstance()->block_compression_mem_tracker());
589
5
        _ctx_pool.clear();
590
5
    }
591
592
265
    Status compress(const Slice& input, faststring* output) override {
593
265
        std::unique_ptr<Context> context;
594
265
        RETURN_IF_ERROR(_acquire_compression_ctx(context));
595
265
        bool compress_failed = false;
596
265
        Defer defer {[&] {
597
265
            if (!compress_failed) {
598
265
                _release_compression_ctx(std::move(context));
599
265
            }
600
265
        }};
601
602
265
        try {
603
265
            Slice compressed_buf;
604
265
            size_t max_len = max_compressed_len(input.size);
605
265
            if (max_len > MAX_COMPRESSION_BUFFER_SIZE_FOR_REUSE) {
606
                // use output directly
607
0
                output->resize(max_len);
608
0
                compressed_buf.data = reinterpret_cast<char*>(output->data());
609
0
                compressed_buf.size = max_len;
610
265
            } else {
611
265
                {
612
265
                    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
613
265
                            ExecEnv::GetInstance()->block_compression_mem_tracker());
614
                    // reuse context buffer if max_len <= MAX_COMPRESSION_BUFFER_FOR_REUSE
615
265
                    context->buffer->resize(max_len);
616
265
                }
617
265
                compressed_buf.data = reinterpret_cast<char*>(context->buffer->data());
618
265
                compressed_buf.size = max_len;
619
265
            }
620
621
265
            size_t compressed_len = LZ4_compress_HC_continue(
622
265
                    context->ctx, input.data, compressed_buf.data, cast_set<int>(input.size),
623
265
                    static_cast<int>(compressed_buf.size));
624
265
            if (compressed_len == 0) {
625
0
                compress_failed = true;
626
0
                return Status::InvalidArgument("Output buffer's capacity is not enough, size={}",
627
0
                                               compressed_buf.size);
628
0
            }
629
265
            output->resize(compressed_len);
630
265
            if (max_len <= MAX_COMPRESSION_BUFFER_SIZE_FOR_REUSE) {
631
265
                output->assign_copy(reinterpret_cast<uint8_t*>(compressed_buf.data),
632
265
                                    compressed_len);
633
265
            }
634
265
        } catch (...) {
635
            // Do not set compress_failed to release context
636
0
            DCHECK(!compress_failed);
637
0
            return Status::InternalError("Fail to do LZ4HC compress due to exception");
638
0
        }
639
265
        return Status::OK();
640
265
    }
641
642
13
    Status decompress(const Slice& input, Slice* output) override {
643
13
        auto decompressed_len = LZ4_decompress_safe(
644
13
                input.data, output->data, cast_set<int>(input.size), cast_set<int>(output->size));
645
13
        if (decompressed_len < 0) {
646
5
            return Status::InvalidArgument(
647
5
                    "destination buffer is not large enough or the source stream is detected "
648
5
                    "malformed, fail to do LZ4 decompress, error={}",
649
5
                    decompressed_len);
650
5
        }
651
8
        output->size = decompressed_len;
652
8
        return Status::OK();
653
13
    }
654
655
265
    size_t max_compressed_len(size_t len) override { return LZ4_compressBound(cast_set<int>(len)); }
656
657
private:
658
265
    Status _acquire_compression_ctx(std::unique_ptr<Context>& out) {
659
265
        std::lock_guard<std::mutex> l(_ctx_mutex);
660
265
        if (_ctx_pool.empty()) {
661
4
            std::unique_ptr<Context> localCtx = Context::create_unique();
662
4
            if (localCtx.get() == nullptr) {
663
0
                return Status::InvalidArgument("new LZ4HC context error");
664
0
            }
665
            // Allocate the native stream under the compression tracker so its
666
            // creation and the destructor's LZ4_freeStreamHC() hit the same tracker.
667
4
            {
668
4
                SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
669
4
                        ExecEnv::GetInstance()->block_compression_mem_tracker());
670
4
                localCtx->ctx = LZ4_createStreamHC();
671
4
            }
672
4
            if (localCtx->ctx == nullptr) {
673
0
                return Status::InvalidArgument("LZ4_createStreamHC error");
674
0
            }
675
            // A newly created stream defaults to the library's default level, so
676
            // apply the requested level here; otherwise the first page compressed
677
            // by this context would ignore the configured level.
678
4
            LZ4_resetStreamHC_fast(localCtx->ctx, static_cast<int>(_compression_level));
679
4
            out = std::move(localCtx);
680
4
            return Status::OK();
681
4
        }
682
261
        out = std::move(_ctx_pool.back());
683
261
        _ctx_pool.pop_back();
684
261
        return Status::OK();
685
265
    }
686
265
    void _release_compression_ctx(std::unique_ptr<Context> context) {
687
265
        DCHECK(context);
688
265
        LZ4_resetStreamHC_fast(context->ctx, static_cast<int>(_compression_level));
689
265
        std::lock_guard<std::mutex> l(_ctx_mutex);
690
265
        _ctx_pool.push_back(std::move(context));
691
265
    }
692
693
private:
694
    int64_t _compression_level = config::LZ4_HC_compression_level;
695
    mutable std::mutex _ctx_mutex;
696
    mutable std::vector<std::unique_ptr<Context>> _ctx_pool;
697
};
698
699
class SnappySlicesSource : public snappy::Source {
700
public:
701
    SnappySlicesSource(const std::vector<Slice>& slices)
702
133
            : _available(0), _cur_slice(0), _slice_off(0) {
703
137
        for (auto& slice : slices) {
704
            // We filter empty slice here to avoid complicated process
705
137
            if (slice.size == 0) {
706
1
                continue;
707
1
            }
708
136
            _available += slice.size;
709
136
            _slices.push_back(slice);
710
136
        }
711
133
    }
712
133
    ~SnappySlicesSource() override {}
713
714
    // Return the number of bytes left to read from the source
715
133
    size_t Available() const override { return _available; }
716
717
    // Peek at the next flat region of the source.  Does not reposition
718
    // the source.  The returned region is empty iff Available()==0.
719
    //
720
    // Returns a pointer to the beginning of the region and store its
721
    // length in *len.
722
    //
723
    // The returned region is valid until the next call to Skip() or
724
    // until this object is destroyed, whichever occurs first.
725
    //
726
    // The returned region may be larger than Available() (for example
727
    // if this ByteSource is a view on a substring of a larger source).
728
    // The caller is responsible for ensuring that it only reads the
729
    // Available() bytes.
730
151
    const char* Peek(size_t* len) override {
731
151
        if (_available == 0) {
732
0
            *len = 0;
733
0
            return nullptr;
734
0
        }
735
        // we should assure that *len is not 0
736
151
        *len = _slices[_cur_slice].size - _slice_off;
737
151
        DCHECK(*len != 0);
738
151
        return _slices[_cur_slice].data + _slice_off;
739
151
    }
740
741
    // Skip the next n bytes.  Invalidates any buffer returned by
742
    // a previous call to Peek().
743
    // REQUIRES: Available() >= n
744
152
    void Skip(size_t n) override {
745
152
        _available -= n;
746
288
        while (n > 0) {
747
151
            auto left = _slices[_cur_slice].size - _slice_off;
748
151
            if (left > n) {
749
                // n can be digest in current slice
750
15
                _slice_off += n;
751
15
                return;
752
15
            }
753
136
            _slice_off = 0;
754
136
            _cur_slice++;
755
136
            n -= left;
756
136
        }
757
152
    }
758
759
private:
760
    std::vector<Slice> _slices;
761
    size_t _available;
762
    size_t _cur_slice;
763
    size_t _slice_off;
764
};
765
766
class SnappyBlockCompression : public BlockCompressionCodec {
767
public:
768
144k
    static SnappyBlockCompression* instance() {
769
144k
        static SnappyBlockCompression s_instance;
770
144k
        return &s_instance;
771
144k
    }
772
    ~SnappyBlockCompression() override = default;
773
774
39.1k
    Status compress(const Slice& input, faststring* output) override {
775
39.1k
        size_t max_len = max_compressed_len(input.size);
776
39.1k
        output->resize(max_len);
777
39.1k
        Slice s(*output);
778
779
39.1k
        snappy::RawCompress(input.data, input.size, s.data, &s.size);
780
39.1k
        output->resize(s.size);
781
39.1k
        return Status::OK();
782
39.1k
    }
783
784
217k
    Status decompress(const Slice& input, Slice* output) override {
785
217k
        size_t uncompressed_size = 0;
786
217k
        if (!snappy::GetUncompressedLength(input.data, input.size, &uncompressed_size)) {
787
0
            return Status::InvalidArgument("Fail to get Snappy uncompressed length");
788
0
        }
789
        // RawUncompress has no capacity argument, so reject an undersized destination first.
790
217k
        if (uncompressed_size > output->size) {
791
0
            return Status::InvalidArgument("Snappy output size {} exceeds buffer capacity {}",
792
0
                                           uncompressed_size, output->size);
793
0
        }
794
217k
        if (!snappy::RawUncompress(input.data, input.size, output->data)) {
795
0
            return Status::InvalidArgument("Fail to do Snappy decompress");
796
0
        }
797
217k
        output->size = uncompressed_size;
798
217k
        return Status::OK();
799
217k
    }
800
801
    Status compress(const std::vector<Slice>& inputs, size_t uncompressed_size,
802
133
                    faststring* output) override {
803
133
        auto max_len = max_compressed_len(uncompressed_size);
804
133
        output->resize(max_len);
805
806
133
        SnappySlicesSource source(inputs);
807
133
        snappy::UncheckedByteArraySink sink(reinterpret_cast<char*>(output->data()));
808
133
        output->resize(snappy::Compress(&source, &sink));
809
133
        return Status::OK();
810
133
    }
811
812
39.3k
    size_t max_compressed_len(size_t len) override { return snappy::MaxCompressedLength(len); }
813
};
814
815
class HadoopSnappyBlockCompression : public SnappyBlockCompression {
816
public:
817
158
    static HadoopSnappyBlockCompression* instance() {
818
158
        static HadoopSnappyBlockCompression s_instance;
819
158
        return &s_instance;
820
158
    }
821
    ~HadoopSnappyBlockCompression() override = default;
822
823
    // hadoop use block compression for snappy
824
    // https://github.com/apache/hadoop/blob/trunk/hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-nativetask/src/main/native/src/codec/SnappyCodec.cc
825
273
    Status compress(const Slice& input, faststring* output) override {
826
        // be same with hadop https://github.com/apache/hadoop/blob/trunk/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/SnappyCodec.java
827
273
        size_t snappy_block_size = config::snappy_compression_block_size;
828
273
        size_t overhead = snappy_block_size / 6 + 32;
829
273
        size_t max_input_size = snappy_block_size - overhead;
830
831
273
        size_t data_len = input.size;
832
273
        char* data = input.data;
833
273
        std::vector<OwnedSlice> buffers;
834
273
        size_t out_len = 0;
835
836
612
        while (data_len > 0) {
837
339
            size_t input_size = std::min(data_len, max_input_size);
838
339
            Slice input_slice(data, input_size);
839
339
            faststring output_data;
840
339
            RETURN_IF_ERROR(SnappyBlockCompression::compress(input_slice, &output_data));
841
339
            out_len += output_data.size();
842
            // the OwnedSlice will be moved here
843
339
            buffers.push_back(output_data.build());
844
339
            data += input_size;
845
339
            data_len -= input_size;
846
339
        }
847
848
        // hadoop block compression: umcompressed_length | compressed_length1 | compressed_data1 | compressed_length2 | compressed_data2 | ...
849
273
        size_t total_output_len = 4 + 4 * buffers.size() + out_len;
850
273
        output->resize(total_output_len);
851
273
        char* output_buffer = (char*)output->data();
852
273
        BigEndian::Store32(output_buffer, cast_set<uint32_t>(input.get_size()));
853
273
        output_buffer += 4;
854
339
        for (const auto& buffer : buffers) {
855
339
            auto slice = buffer.slice();
856
339
            BigEndian::Store32(output_buffer, cast_set<uint32_t>(slice.get_size()));
857
339
            output_buffer += 4;
858
339
            memcpy(output_buffer, slice.get_data(), slice.get_size());
859
339
            output_buffer += slice.get_size();
860
339
        }
861
862
273
        DCHECK_EQ(output_buffer - (char*)output->data(), total_output_len);
863
864
273
        return Status::OK();
865
273
    }
866
867
0
    Status decompress(const Slice& input, Slice* output) override {
868
0
        return Status::InternalError("unimplement: SnappyHadoopBlockCompression::decompress");
869
0
    }
870
};
871
872
class ZlibBlockCompression : public BlockCompressionCodec {
873
public:
874
50
    static ZlibBlockCompression* instance() {
875
50
        static ZlibBlockCompression s_instance;
876
50
        return &s_instance;
877
50
    }
878
    ~ZlibBlockCompression() override = default;
879
880
21
    Status compress(const Slice& input, faststring* output) override {
881
21
        size_t max_len = max_compressed_len(input.size);
882
21
        output->resize(max_len);
883
21
        Slice s(*output);
884
885
21
        auto zres = ::compress((Bytef*)s.data, &s.size, (Bytef*)input.data, input.size);
886
21
        if (zres == Z_MEM_ERROR) {
887
0
            throw Exception(Status::MemoryLimitExceeded(fmt::format(
888
0
                    "ZLib compression failed due to memory allocationerror.error = {}, res = {} ",
889
0
                    zError(zres), zres)));
890
21
        } else if (zres != Z_OK) {
891
0
            return Status::InternalError("Fail to do Zlib compress, error={}", zError(zres));
892
0
        }
893
21
        output->resize(s.size);
894
21
        return Status::OK();
895
21
    }
896
897
    Status compress(const std::vector<Slice>& inputs, size_t uncompressed_size,
898
129
                    faststring* output) override {
899
129
        size_t max_len = max_compressed_len(uncompressed_size);
900
129
        output->resize(max_len);
901
902
129
        z_stream zstrm;
903
129
        zstrm.zalloc = Z_NULL;
904
129
        zstrm.zfree = Z_NULL;
905
129
        zstrm.opaque = Z_NULL;
906
129
        auto zres = deflateInit(&zstrm, Z_DEFAULT_COMPRESSION);
907
129
        if (zres == Z_MEM_ERROR) {
908
0
            throw Exception(Status::MemoryLimitExceeded(
909
0
                    "Fail to do ZLib stream compress, error={}, res={}", zError(zres), zres));
910
129
        } else if (zres != Z_OK) {
911
0
            return Status::InternalError("Fail to do ZLib stream compress, error={}, res={}",
912
0
                                         zError(zres), zres);
913
0
        }
914
        // we assume that output is e
915
129
        zstrm.next_out = (Bytef*)output->data();
916
129
        zstrm.avail_out = cast_set<decltype(zstrm.avail_out)>(output->size());
917
262
        for (int i = 0; i < inputs.size(); ++i) {
918
133
            if (inputs[i].size == 0) {
919
1
                continue;
920
1
            }
921
132
            zstrm.next_in = (Bytef*)inputs[i].data;
922
132
            zstrm.avail_in = cast_set<decltype(zstrm.avail_in)>(inputs[i].size);
923
132
            int flush = (i == (inputs.size() - 1)) ? Z_FINISH : Z_NO_FLUSH;
924
925
132
            zres = deflate(&zstrm, flush);
926
132
            if (zres != Z_OK && zres != Z_STREAM_END) {
927
0
                return Status::InternalError("Fail to do ZLib stream compress, error={}, res={}",
928
0
                                             zError(zres), zres);
929
0
            }
930
132
        }
931
932
129
        output->resize(zstrm.total_out);
933
129
        zres = deflateEnd(&zstrm);
934
129
        if (zres == Z_DATA_ERROR) {
935
0
            return Status::InvalidArgument("Fail to do deflateEnd, error={}, res={}", zError(zres),
936
0
                                           zres);
937
129
        } else if (zres != Z_OK) {
938
0
            return Status::InternalError("Fail to do deflateEnd on ZLib stream, error={}, res={}",
939
0
                                         zError(zres), zres);
940
0
        }
941
129
        return Status::OK();
942
129
    }
943
944
162
    Status decompress(const Slice& input, Slice* output) override {
945
162
        size_t input_size = input.size;
946
162
        auto zres =
947
162
                ::uncompress2((Bytef*)output->data, &output->size, (Bytef*)input.data, &input_size);
948
162
        if (zres == Z_DATA_ERROR) {
949
4
            return Status::InvalidArgument("Fail to do ZLib decompress, error={}", zError(zres));
950
158
        } else if (zres == Z_MEM_ERROR) {
951
0
            throw Exception(Status::MemoryLimitExceeded("Fail to do ZLib decompress, error={}",
952
0
                                                        zError(zres)));
953
158
        } else if (zres != Z_OK) {
954
7
            return Status::InternalError("Fail to do ZLib decompress, error={}", zError(zres));
955
7
        }
956
151
        return Status::OK();
957
162
    }
958
959
150
    size_t max_compressed_len(size_t len) override {
960
        // one-time overhead of six bytes for the entire stream plus five bytes per 16 KB block
961
150
        return len + 6 + 5 * ((len >> 14) + 1);
962
150
    }
963
};
964
965
class Bzip2BlockCompression : public BlockCompressionCodec {
966
public:
967
138
    static Bzip2BlockCompression* instance() {
968
138
        static Bzip2BlockCompression s_instance;
969
138
        return &s_instance;
970
138
    }
971
    ~Bzip2BlockCompression() override = default;
972
973
281
    Status compress(const Slice& input, faststring* output) override {
974
281
        size_t max_len = max_compressed_len(input.size);
975
281
        output->resize(max_len);
976
281
        auto size = cast_set<uint32_t>(output->size());
977
281
        auto bzres = BZ2_bzBuffToBuffCompress((char*)output->data(), &size, (char*)input.data,
978
281
                                              cast_set<uint32_t>(input.size), 9, 0, 0);
979
281
        if (bzres == BZ_MEM_ERROR) {
980
0
            throw Exception(
981
0
                    Status::MemoryLimitExceeded("Fail to do Bzip2 compress, ret={}", bzres));
982
281
        } else if (bzres == BZ_PARAM_ERROR) {
983
0
            return Status::InvalidArgument("Fail to do Bzip2 compress, ret={}", bzres);
984
281
        } else if (bzres != BZ_RUN_OK && bzres != BZ_FLUSH_OK && bzres != BZ_FINISH_OK &&
985
281
                   bzres != BZ_STREAM_END && bzres != BZ_OK) {
986
0
            return Status::InternalError("Failed to init bz2. status code: {}", bzres);
987
0
        }
988
281
        output->resize(size);
989
281
        return Status::OK();
990
281
    }
991
992
    Status compress(const std::vector<Slice>& inputs, size_t uncompressed_size,
993
0
                    faststring* output) override {
994
0
        size_t max_len = max_compressed_len(uncompressed_size);
995
0
        output->resize(max_len);
996
997
0
        bz_stream bzstrm;
998
0
        bzero(&bzstrm, sizeof(bzstrm));
999
0
        int bzres = BZ2_bzCompressInit(&bzstrm, 9, 0, 0);
1000
0
        if (bzres == BZ_PARAM_ERROR) {
1001
0
            return Status::InvalidArgument("Failed to init bz2. status code: {}", bzres);
1002
0
        } else if (bzres == BZ_MEM_ERROR) {
1003
0
            throw Exception(
1004
0
                    Status::MemoryLimitExceeded("Failed to init bz2. status code: {}", bzres));
1005
0
        } else if (bzres != BZ_OK) {
1006
0
            return Status::InternalError("Failed to init bz2. status code: {}", bzres);
1007
0
        }
1008
        // we assume that output is e
1009
0
        bzstrm.next_out = (char*)output->data();
1010
0
        bzstrm.avail_out = cast_set<uint32_t>(output->size());
1011
0
        for (int i = 0; i < inputs.size(); ++i) {
1012
0
            if (inputs[i].size == 0) {
1013
0
                continue;
1014
0
            }
1015
0
            bzstrm.next_in = (char*)inputs[i].data;
1016
0
            bzstrm.avail_in = cast_set<uint32_t>(inputs[i].size);
1017
0
            int flush = (i == (inputs.size() - 1)) ? BZ_FINISH : BZ_RUN;
1018
1019
0
            bzres = BZ2_bzCompress(&bzstrm, flush);
1020
0
            if (bzres == BZ_PARAM_ERROR) {
1021
0
                return Status::InvalidArgument("Failed to init bz2. status code: {}", bzres);
1022
0
            } else if (bzres != BZ_RUN_OK && bzres != BZ_FLUSH_OK && bzres != BZ_FINISH_OK &&
1023
0
                       bzres != BZ_STREAM_END && bzres != BZ_OK) {
1024
0
                return Status::InternalError("Failed to init bz2. status code: {}", bzres);
1025
0
            }
1026
0
        }
1027
1028
0
        size_t total_out = (size_t)bzstrm.total_out_hi32 << 32 | (size_t)bzstrm.total_out_lo32;
1029
0
        output->resize(total_out);
1030
0
        bzres = BZ2_bzCompressEnd(&bzstrm);
1031
0
        if (bzres == BZ_PARAM_ERROR) {
1032
0
            return Status::InvalidArgument("Fail to do deflateEnd on bzip2 stream, res={}", bzres);
1033
0
        } else if (bzres != BZ_OK) {
1034
0
            return Status::InternalError("Fail to do deflateEnd on bzip2 stream, res={}", bzres);
1035
0
        }
1036
0
        return Status::OK();
1037
0
    }
1038
1039
0
    Status decompress(const Slice& input, Slice* output) override {
1040
0
        return Status::InternalError("unimplement: Bzip2BlockCompression::decompress");
1041
0
    }
1042
1043
281
    size_t max_compressed_len(size_t len) override {
1044
        // TODO: make sure the max_compressed_len for bzip2
1045
        // 50 is an estimate fix overhead for bzip2
1046
        // in case the input len is small and BZ2_bzBuffToBuffCompress will return
1047
        // BZ_OUTBUFF_FULL
1048
281
        return len * 2 + 50;
1049
281
    }
1050
};
1051
1052
// for ZSTD compression and decompression, with BOTH fast and high compression ratio
1053
class ZstdBlockCompression : public BlockCompressionCodec {
1054
private:
1055
    class CContext {
1056
        ENABLE_FACTORY_CREATOR(CContext);
1057
1058
    public:
1059
43
        CContext() : ctx(nullptr) {
1060
43
            SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
1061
43
                    ExecEnv::GetInstance()->block_compression_mem_tracker());
1062
43
            buffer = std::make_unique<faststring>();
1063
43
        }
1064
        ZSTD_CCtx* ctx;
1065
        std::unique_ptr<faststring> buffer;
1066
12
        ~CContext() {
1067
12
            SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
1068
12
                    ExecEnv::GetInstance()->block_compression_mem_tracker());
1069
12
            if (ctx) {
1070
12
                ZSTD_freeCCtx(ctx);
1071
12
            }
1072
12
            buffer.reset();
1073
12
        }
1074
    };
1075
    class DContext {
1076
        ENABLE_FACTORY_CREATOR(DContext);
1077
1078
    public:
1079
37
        DContext() : ctx(nullptr) {}
1080
        ZSTD_DCtx* ctx;
1081
20
        ~DContext() {
1082
20
            if (ctx) {
1083
20
                ZSTD_freeDCtx(ctx);
1084
20
            }
1085
20
        }
1086
    };
1087
1088
public:
1089
22.2M
    static ZstdBlockCompression* instance() {
1090
22.2M
        static ZstdBlockCompression s_instance;
1091
22.2M
        return &s_instance;
1092
22.2M
    }
1093
    ZstdBlockCompression() = default;
1094
5
    explicit ZstdBlockCompression(int level) : _compression_level(level) {}
1095
9
    ~ZstdBlockCompression() {
1096
9
        SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
1097
9
                ExecEnv::GetInstance()->block_compression_mem_tracker());
1098
9
        _ctx_c_pool.clear();
1099
9
        _ctx_d_pool.clear();
1100
9
    }
1101
1102
369k
    size_t max_compressed_len(size_t len) override { return ZSTD_compressBound(len); }
1103
1104
978
    Status compress(const Slice& input, faststring* output) override {
1105
978
        std::vector<Slice> inputs {input};
1106
978
        return compress(inputs, input.size, output);
1107
978
    }
1108
1109
    // follow ZSTD official example
1110
    //  https://github.com/facebook/zstd/blob/dev/examples/streaming_compression.c
1111
    Status compress(const std::vector<Slice>& inputs, size_t uncompressed_size,
1112
369k
                    faststring* output) override {
1113
369k
        std::unique_ptr<CContext> context;
1114
369k
        RETURN_IF_ERROR(_acquire_compression_ctx(context));
1115
369k
        bool compress_failed = false;
1116
369k
        Defer defer {[&] {
1117
369k
            if (!compress_failed) {
1118
369k
                _release_compression_ctx(std::move(context));
1119
369k
            }
1120
369k
        }};
1121
1122
369k
        try {
1123
369k
            size_t max_len = max_compressed_len(uncompressed_size);
1124
369k
            Slice compressed_buf;
1125
369k
            if (max_len > MAX_COMPRESSION_BUFFER_SIZE_FOR_REUSE) {
1126
                // use output directly
1127
4
                output->resize(max_len);
1128
4
                compressed_buf.data = reinterpret_cast<char*>(output->data());
1129
4
                compressed_buf.size = max_len;
1130
369k
            } else {
1131
369k
                {
1132
369k
                    SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
1133
369k
                            ExecEnv::GetInstance()->block_compression_mem_tracker());
1134
                    // reuse context buffer if max_len <= MAX_COMPRESSION_BUFFER_FOR_REUSE
1135
369k
                    context->buffer->resize(max_len);
1136
369k
                }
1137
369k
                compressed_buf.data = reinterpret_cast<char*>(context->buffer->data());
1138
369k
                compressed_buf.size = max_len;
1139
369k
            }
1140
1141
369k
            auto ret = ZSTD_CCtx_setParameter(context->ctx, ZSTD_c_compressionLevel,
1142
369k
                                              _compression_level);
1143
369k
            if (ZSTD_isError(ret)) {
1144
0
                return Status::InvalidArgument("ZSTD_CCtx_setParameter compression level error: {}",
1145
0
                                               ZSTD_getErrorString(ZSTD_getErrorCode(ret)));
1146
0
            }
1147
            // set checksum flag to 1
1148
369k
            ret = ZSTD_CCtx_setParameter(context->ctx, ZSTD_c_checksumFlag, 1);
1149
369k
            if (ZSTD_isError(ret)) {
1150
0
                return Status::InvalidArgument("ZSTD_CCtx_setParameter checksumFlag error: {}",
1151
0
                                               ZSTD_getErrorString(ZSTD_getErrorCode(ret)));
1152
0
            }
1153
1154
369k
            ZSTD_outBuffer out_buf = {compressed_buf.data, compressed_buf.size, 0};
1155
1156
776k
            for (size_t i = 0; i < inputs.size(); i++) {
1157
407k
                ZSTD_inBuffer in_buf = {inputs[i].data, inputs[i].size, 0};
1158
1159
407k
                bool last_input = (i == inputs.size() - 1);
1160
407k
                auto mode = last_input ? ZSTD_e_end : ZSTD_e_continue;
1161
1162
407k
                bool finished = false;
1163
407k
                do {
1164
                    // do compress
1165
407k
                    ret = ZSTD_compressStream2(context->ctx, &out_buf, &in_buf, mode);
1166
1167
407k
                    if (ZSTD_isError(ret)) {
1168
0
                        compress_failed = true;
1169
0
                        return Status::InternalError("ZSTD_compressStream2 error: {}",
1170
0
                                                     ZSTD_getErrorString(ZSTD_getErrorCode(ret)));
1171
0
                    }
1172
1173
                    // ret is ZSTD hint for needed output buffer size
1174
407k
                    if (ret > 0 && out_buf.pos == out_buf.size) {
1175
0
                        compress_failed = true;
1176
0
                        return Status::InternalError("ZSTD_compressStream2 output buffer full");
1177
0
                    }
1178
1179
407k
                    finished = last_input ? (ret == 0) : (in_buf.pos == inputs[i].size);
1180
407k
                } while (!finished);
1181
407k
            }
1182
1183
            // set compressed size for caller
1184
369k
            output->resize(out_buf.pos);
1185
369k
            if (max_len <= MAX_COMPRESSION_BUFFER_SIZE_FOR_REUSE) {
1186
369k
                output->assign_copy(reinterpret_cast<uint8_t*>(compressed_buf.data), out_buf.pos);
1187
369k
            }
1188
369k
        } catch (std::exception& e) {
1189
0
            return Status::InternalError("Fail to do ZSTD compress due to exception {}", e.what());
1190
0
        } catch (...) {
1191
            // Do not set compress_failed to release context
1192
0
            DCHECK(!compress_failed);
1193
0
            return Status::InternalError("Fail to do ZSTD compress due to exception");
1194
0
        }
1195
1196
369k
        return Status::OK();
1197
369k
    }
1198
1199
209k
    Status decompress(const Slice& input, Slice* output) override {
1200
209k
        std::unique_ptr<DContext> context;
1201
209k
        bool decompress_failed = false;
1202
209k
        RETURN_IF_ERROR(_acquire_decompression_ctx(context));
1203
209k
        Defer defer {[&] {
1204
209k
            if (!decompress_failed) {
1205
209k
                _release_decompression_ctx(std::move(context));
1206
209k
            }
1207
209k
        }};
1208
1209
209k
        size_t ret = ZSTD_decompressDCtx(context->ctx, output->data, output->size, input.data,
1210
209k
                                         input.size);
1211
209k
        if (ZSTD_isError(ret)) {
1212
5
            decompress_failed = true;
1213
5
            return Status::InternalError("ZSTD_decompressDCtx error: {}",
1214
5
                                         ZSTD_getErrorString(ZSTD_getErrorCode(ret)));
1215
5
        }
1216
1217
        // set decompressed size for caller
1218
209k
        output->size = ret;
1219
1220
209k
        return Status::OK();
1221
209k
    }
1222
1223
private:
1224
369k
    Status _acquire_compression_ctx(std::unique_ptr<CContext>& out) {
1225
369k
        std::lock_guard<std::mutex> l(_ctx_c_mutex);
1226
369k
        if (_ctx_c_pool.empty()) {
1227
43
            std::unique_ptr<CContext> localCtx = CContext::create_unique();
1228
43
            if (localCtx.get() == nullptr) {
1229
0
                return Status::InvalidArgument("failed to new ZSTD CContext");
1230
0
            }
1231
            //typedef LZ4F_cctx* LZ4F_compressionContext_t;
1232
            // Allocate the native context under the compression tracker so its
1233
            // creation and the destructor's ZSTD_freeCCtx() hit the same tracker.
1234
43
            {
1235
43
                SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
1236
43
                        ExecEnv::GetInstance()->block_compression_mem_tracker());
1237
43
                localCtx->ctx = ZSTD_createCCtx();
1238
43
            }
1239
43
            if (localCtx->ctx == nullptr) {
1240
0
                return Status::InvalidArgument("Failed to create ZSTD compress ctx");
1241
0
            }
1242
43
            out = std::move(localCtx);
1243
43
            return Status::OK();
1244
43
        }
1245
368k
        out = std::move(_ctx_c_pool.back());
1246
368k
        _ctx_c_pool.pop_back();
1247
368k
        return Status::OK();
1248
369k
    }
1249
369k
    void _release_compression_ctx(std::unique_ptr<CContext> context) {
1250
369k
        DCHECK(context);
1251
369k
        auto ret = ZSTD_CCtx_reset(context->ctx, ZSTD_reset_session_only);
1252
369k
        DCHECK(!ZSTD_isError(ret));
1253
369k
        std::lock_guard<std::mutex> l(_ctx_c_mutex);
1254
369k
        _ctx_c_pool.push_back(std::move(context));
1255
369k
    }
1256
1257
209k
    Status _acquire_decompression_ctx(std::unique_ptr<DContext>& out) {
1258
209k
        std::lock_guard<std::mutex> l(_ctx_d_mutex);
1259
209k
        if (_ctx_d_pool.empty()) {
1260
37
            std::unique_ptr<DContext> localCtx = DContext::create_unique();
1261
37
            if (localCtx.get() == nullptr) {
1262
0
                return Status::InvalidArgument("failed to new ZSTD DContext");
1263
0
            }
1264
37
            localCtx->ctx = ZSTD_createDCtx();
1265
37
            if (localCtx->ctx == nullptr) {
1266
0
                return Status::InvalidArgument("Fail to init ZSTD decompress context");
1267
0
            }
1268
37
            out = std::move(localCtx);
1269
37
            return Status::OK();
1270
37
        }
1271
209k
        out = std::move(_ctx_d_pool.back());
1272
209k
        _ctx_d_pool.pop_back();
1273
209k
        return Status::OK();
1274
209k
    }
1275
209k
    void _release_decompression_ctx(std::unique_ptr<DContext> context) {
1276
209k
        DCHECK(context);
1277
        // reset ctx to start a new decompress session
1278
209k
        auto ret = ZSTD_DCtx_reset(context->ctx, ZSTD_reset_session_only);
1279
209k
        DCHECK(!ZSTD_isError(ret));
1280
209k
        std::lock_guard<std::mutex> l(_ctx_d_mutex);
1281
209k
        _ctx_d_pool.push_back(std::move(context));
1282
209k
    }
1283
1284
private:
1285
    int _compression_level = ZSTD_CLEVEL_DEFAULT;
1286
    mutable std::mutex _ctx_c_mutex;
1287
    mutable std::vector<std::unique_ptr<CContext>> _ctx_c_pool;
1288
1289
    mutable std::mutex _ctx_d_mutex;
1290
    mutable std::vector<std::unique_ptr<DContext>> _ctx_d_pool;
1291
};
1292
1293
class GzipBlockCompression : public ZlibBlockCompression {
1294
public:
1295
153
    static GzipBlockCompression* instance() {
1296
153
        static GzipBlockCompression s_instance;
1297
153
        return &s_instance;
1298
153
    }
1299
    ~GzipBlockCompression() override = default;
1300
1301
324
    Status compress(const Slice& input, faststring* output) override {
1302
324
        size_t max_len = max_compressed_len(input.size);
1303
324
        output->resize(max_len);
1304
1305
324
        z_stream z_strm = {};
1306
324
        z_strm.zalloc = Z_NULL;
1307
324
        z_strm.zfree = Z_NULL;
1308
324
        z_strm.opaque = Z_NULL;
1309
1310
324
        int zres = deflateInit2(&z_strm, Z_DEFAULT_COMPRESSION, Z_DEFLATED, MAX_WBITS + GZIP_CODEC,
1311
324
                                8, Z_DEFAULT_STRATEGY);
1312
1313
324
        if (zres == Z_MEM_ERROR) {
1314
0
            throw Exception(Status::MemoryLimitExceeded(
1315
0
                    "Fail to init ZLib compress, error={}, res={}", zError(zres), zres));
1316
324
        } else if (zres != Z_OK) {
1317
0
            return Status::InternalError("Fail to init ZLib compress, error={}, res={}",
1318
0
                                         zError(zres), zres);
1319
0
        }
1320
1321
324
        z_strm.next_in = (Bytef*)input.get_data();
1322
324
        z_strm.avail_in = cast_set<decltype(z_strm.avail_in)>(input.get_size());
1323
324
        z_strm.next_out = (Bytef*)output->data();
1324
324
        z_strm.avail_out = cast_set<decltype(z_strm.avail_out)>(output->size());
1325
1326
324
        zres = deflate(&z_strm, Z_FINISH);
1327
324
        if (zres != Z_OK && zres != Z_STREAM_END) {
1328
0
            return Status::InternalError("Fail to do ZLib stream compress, error={}, res={}",
1329
0
                                         zError(zres), zres);
1330
0
        }
1331
1332
324
        output->resize(z_strm.total_out);
1333
324
        zres = deflateEnd(&z_strm);
1334
324
        if (zres == Z_DATA_ERROR) {
1335
0
            return Status::InvalidArgument("Fail to end zlib compress");
1336
324
        } else if (zres != Z_OK) {
1337
0
            return Status::InternalError("Fail to end zlib compress");
1338
0
        }
1339
324
        return Status::OK();
1340
324
    }
1341
1342
    Status compress(const std::vector<Slice>& inputs, size_t uncompressed_size,
1343
0
                    faststring* output) override {
1344
0
        size_t max_len = max_compressed_len(uncompressed_size);
1345
0
        output->resize(max_len);
1346
1347
0
        z_stream zstrm;
1348
0
        zstrm.zalloc = Z_NULL;
1349
0
        zstrm.zfree = Z_NULL;
1350
0
        zstrm.opaque = Z_NULL;
1351
0
        auto zres = deflateInit2(&zstrm, Z_DEFAULT_COMPRESSION, Z_DEFLATED, MAX_WBITS + GZIP_CODEC,
1352
0
                                 8, Z_DEFAULT_STRATEGY);
1353
0
        if (zres == Z_MEM_ERROR) {
1354
0
            throw Exception(Status::MemoryLimitExceeded(
1355
0
                    "Fail to init ZLib stream compress, error={}, res={}", zError(zres), zres));
1356
0
        } else if (zres != Z_OK) {
1357
0
            return Status::InternalError("Fail to init ZLib stream compress, error={}, res={}",
1358
0
                                         zError(zres), zres);
1359
0
        }
1360
1361
        // we assume that output is e
1362
0
        zstrm.next_out = (Bytef*)output->data();
1363
0
        zstrm.avail_out = cast_set<decltype(zstrm.avail_out)>(output->size());
1364
0
        for (int i = 0; i < inputs.size(); ++i) {
1365
0
            if (inputs[i].size == 0) {
1366
0
                continue;
1367
0
            }
1368
0
            zstrm.next_in = (Bytef*)inputs[i].data;
1369
0
            zstrm.avail_in = cast_set<decltype(zstrm.avail_in)>(inputs[i].size);
1370
0
            int flush = (i == (inputs.size() - 1)) ? Z_FINISH : Z_NO_FLUSH;
1371
1372
0
            zres = deflate(&zstrm, flush);
1373
0
            if (zres != Z_OK && zres != Z_STREAM_END) {
1374
0
                return Status::InternalError("Fail to do ZLib stream compress, error={}, res={}",
1375
0
                                             zError(zres), zres);
1376
0
            }
1377
0
        }
1378
1379
0
        output->resize(zstrm.total_out);
1380
0
        zres = deflateEnd(&zstrm);
1381
0
        if (zres == Z_DATA_ERROR) {
1382
0
            return Status::InvalidArgument("Fail to do deflateEnd on ZLib stream, error={}, res={}",
1383
0
                                           zError(zres), zres);
1384
0
        } else if (zres != Z_OK) {
1385
0
            return Status::InternalError("Fail to do deflateEnd on ZLib stream, error={}, res={}",
1386
0
                                         zError(zres), zres);
1387
0
        }
1388
0
        return Status::OK();
1389
0
    }
1390
1391
0
    Status decompress(const Slice& input, Slice* output) override {
1392
0
        z_stream z_strm = {};
1393
0
        z_strm.zalloc = Z_NULL;
1394
0
        z_strm.zfree = Z_NULL;
1395
0
        z_strm.opaque = Z_NULL;
1396
1397
0
        int ret = inflateInit2(&z_strm, MAX_WBITS + GZIP_CODEC);
1398
0
        if (ret != Z_OK) {
1399
0
            return Status::InternalError("Fail to init ZLib decompress, error={}, res={}",
1400
0
                                         zError(ret), ret);
1401
0
        }
1402
1403
        // 1. set input and output
1404
0
        z_strm.next_in = reinterpret_cast<Bytef*>(input.data);
1405
0
        z_strm.avail_in = cast_set<decltype(z_strm.avail_in)>(input.size);
1406
0
        z_strm.next_out = reinterpret_cast<Bytef*>(output->data);
1407
0
        z_strm.avail_out = cast_set<decltype(z_strm.avail_out)>(output->size);
1408
1409
0
        if (z_strm.avail_out > 0) {
1410
            // We only support non-streaming use case  for block decompressor
1411
0
            ret = inflate(&z_strm, Z_FINISH);
1412
0
            if (ret != Z_OK && ret != Z_STREAM_END) {
1413
0
                (void)inflateEnd(&z_strm);
1414
0
                if (ret == Z_MEM_ERROR) {
1415
0
                    throw Exception(Status::MemoryLimitExceeded(
1416
0
                            "Fail to do ZLib stream compress, error={}, res={}", zError(ret), ret));
1417
0
                } else if (ret == Z_DATA_ERROR) {
1418
0
                    return Status::InvalidArgument(
1419
0
                            "Fail to do ZLib stream compress, error={}, res={}", zError(ret), ret);
1420
0
                }
1421
0
                return Status::InternalError("Fail to do ZLib stream compress, error={}, res={}",
1422
0
                                             zError(ret), ret);
1423
0
            }
1424
0
        }
1425
0
        (void)inflateEnd(&z_strm);
1426
1427
0
        return Status::OK();
1428
0
    }
1429
1430
324
    size_t max_compressed_len(size_t len) override {
1431
324
        z_stream zstrm;
1432
324
        zstrm.zalloc = Z_NULL;
1433
324
        zstrm.zfree = Z_NULL;
1434
324
        zstrm.opaque = Z_NULL;
1435
324
        auto zres = deflateInit2(&zstrm, Z_DEFAULT_COMPRESSION, Z_DEFLATED, MAX_WBITS + GZIP_CODEC,
1436
324
                                 MEM_LEVEL, Z_DEFAULT_STRATEGY);
1437
324
        if (zres != Z_OK) {
1438
            // Fall back to zlib estimate logic for deflate, notice this may
1439
            // cause decompress error
1440
0
            LOG(WARNING) << "Fail to do ZLib stream compress, error=" << zError(zres)
1441
0
                         << ", res=" << zres;
1442
0
            return ZlibBlockCompression::max_compressed_len(len);
1443
324
        } else {
1444
324
            zres = deflateEnd(&zstrm);
1445
324
            if (zres != Z_OK) {
1446
0
                LOG(WARNING) << "Fail to do deflateEnd on ZLib stream, error=" << zError(zres)
1447
0
                             << ", res=" << zres;
1448
0
            }
1449
            // Mark, maintainer of zlib, has stated that 12 needs to be added to
1450
            // result for gzip
1451
            // http://compgroups.net/comp.unix.programmer/gzip-compressing-an-in-memory-string-usi/54854
1452
            // To have a safe upper bound for "wrapper variations", we add 32 to
1453
            // estimate
1454
324
            auto upper_bound = deflateBound(&zstrm, len) + 32;
1455
324
            return upper_bound;
1456
324
        }
1457
324
    }
1458
1459
private:
1460
    // Magic number for zlib, see https://zlib.net/manual.html for more details.
1461
    const static int GZIP_CODEC = 16; // gzip
1462
    // The memLevel parameter specifies how much memory should be allocated for
1463
    // the internal compression state.
1464
    const static int MEM_LEVEL = 8;
1465
};
1466
1467
// Only used on x86 or x86_64
1468
#if defined(__x86_64__) || defined(_M_X64) || defined(i386) || defined(__i386__) || \
1469
        defined(__i386) || defined(_M_IX86)
1470
class GzipBlockCompressionByLibdeflate final : public GzipBlockCompression {
1471
public:
1472
3
    GzipBlockCompressionByLibdeflate() : GzipBlockCompression() {}
1473
10.0k
    static GzipBlockCompressionByLibdeflate* instance() {
1474
10.0k
        static GzipBlockCompressionByLibdeflate s_instance;
1475
10.0k
        return &s_instance;
1476
10.0k
    }
1477
    ~GzipBlockCompressionByLibdeflate() override = default;
1478
1479
11.7k
    Status decompress(const Slice& input, Slice* output) override {
1480
11.7k
        if (input.empty()) {
1481
0
            output->size = 0;
1482
0
            return Status::OK();
1483
0
        }
1484
11.7k
        thread_local std::unique_ptr<libdeflate_decompressor, void (*)(libdeflate_decompressor*)>
1485
11.7k
                decompressor {libdeflate_alloc_decompressor(), libdeflate_free_decompressor};
1486
11.7k
        if (!decompressor) {
1487
0
            return Status::InternalError("libdeflate_alloc_decompressor error.");
1488
0
        }
1489
11.7k
        std::size_t out_len;
1490
11.7k
        auto result = libdeflate_gzip_decompress(decompressor.get(), input.data, input.size,
1491
11.7k
                                                 output->data, output->size, &out_len);
1492
11.7k
        if (result != LIBDEFLATE_SUCCESS) {
1493
0
            return Status::InternalError("libdeflate_gzip_decompress error, res={}", result);
1494
0
        }
1495
11.7k
        return Status::OK();
1496
11.7k
    }
1497
};
1498
#endif
1499
1500
class LzoBlockCompression final : public BlockCompressionCodec {
1501
public:
1502
418
    static LzoBlockCompression* instance() {
1503
418
        static LzoBlockCompression s_instance;
1504
418
        return &s_instance;
1505
418
    }
1506
1507
0
    Status compress(const Slice& input, faststring* output) override {
1508
0
        return Status::InvalidArgument("not impl lzo compress.");
1509
0
    }
1510
0
    size_t max_compressed_len(size_t len) override { return 0; };
1511
430
    Status decompress(const Slice& input, Slice* output) override {
1512
430
        auto* input_ptr = input.data;
1513
430
        auto remain_input_size = input.size;
1514
430
        auto* output_ptr = output->data;
1515
430
        auto remain_output_size = output->size;
1516
430
        auto* output_limit = output->data + output->size;
1517
1518
        // Example:
1519
        // OriginData(The original data will be divided into several large data block.) :
1520
        //      large data block1 | large data block2 | large data block3 | ....
1521
        // The large data block will be divided into several small data block.
1522
        // Suppose a large data block is divided into three small blocks:
1523
        // large data block1:            | small block1 | small block2 | small block3 |
1524
        // CompressData:   <A [B1 compress(small block1) ] [B2 compress(small block1) ] [B3 compress(small block1)]>
1525
        //
1526
        // A : original length of the current block of large data block.
1527
        // sizeof(A) = 4 bytes.
1528
        // A = length(small block1) + length(small block2) + length(small block3)
1529
        // Bx : length of  small data block bx.
1530
        // sizeof(Bx) = 4 bytes.
1531
        // Bx = length(compress(small blockx))
1532
430
        try {
1533
860
            while (remain_input_size > 0) {
1534
430
                if (remain_input_size < 4) {
1535
0
                    return Status::InvalidArgument(
1536
0
                            "Need more input buffer to get large_block_uncompressed_len.");
1537
0
                }
1538
1539
430
                uint32_t large_block_uncompressed_len = BigEndian::Load32(input_ptr);
1540
430
                input_ptr += 4;
1541
430
                remain_input_size -= 4;
1542
1543
430
                if (remain_output_size < large_block_uncompressed_len) {
1544
0
                    return Status::InvalidArgument(
1545
0
                            "Need more output buffer to get uncompressed data.");
1546
0
                }
1547
1548
860
                while (large_block_uncompressed_len > 0) {
1549
430
                    if (remain_input_size < 4) {
1550
0
                        return Status::InvalidArgument(
1551
0
                                "Need more input buffer to get small_block_compressed_len.");
1552
0
                    }
1553
1554
430
                    uint32_t small_block_compressed_len = BigEndian::Load32(input_ptr);
1555
430
                    input_ptr += 4;
1556
430
                    remain_input_size -= 4;
1557
1558
430
                    if (remain_input_size < small_block_compressed_len) {
1559
0
                        return Status::InvalidArgument(
1560
0
                                "Need more input buffer to decompress small block.");
1561
0
                    }
1562
1563
430
                    auto small_block_uncompressed_len =
1564
430
                            orc::lzoDecompress(input_ptr, input_ptr + small_block_compressed_len,
1565
430
                                               output_ptr, output_limit);
1566
1567
430
                    input_ptr += small_block_compressed_len;
1568
430
                    remain_input_size -= small_block_compressed_len;
1569
1570
430
                    output_ptr += small_block_uncompressed_len;
1571
430
                    large_block_uncompressed_len -= small_block_uncompressed_len;
1572
430
                    remain_output_size -= small_block_uncompressed_len;
1573
430
                }
1574
430
            }
1575
430
        } catch (const orc::ParseError& e) {
1576
            //Prevent be from hanging due to orc::lzoDecompress throw exception
1577
0
            return Status::InternalError("Fail to do LZO decompress, error={}", e.what());
1578
0
        }
1579
430
        return Status::OK();
1580
430
    }
1581
};
1582
1583
class BrotliBlockCompression final : public BlockCompressionCodec {
1584
public:
1585
32
    static BrotliBlockCompression* instance() {
1586
32
        static BrotliBlockCompression s_instance;
1587
32
        return &s_instance;
1588
32
    }
1589
1590
0
    Status compress(const Slice& input, faststring* output) override {
1591
0
        return Status::InvalidArgument("not impl brotli compress.");
1592
0
    }
1593
1594
0
    size_t max_compressed_len(size_t len) override { return 0; };
1595
1596
66
    Status decompress(const Slice& input, Slice* output) override {
1597
        // The size of output buffer is always equal to the umcompressed length.
1598
66
        BrotliDecoderResult result = BrotliDecoderDecompress(
1599
66
                input.get_size(), reinterpret_cast<const uint8_t*>(input.get_data()), &output->size,
1600
66
                reinterpret_cast<uint8_t*>(output->data));
1601
66
        if (result != BROTLI_DECODER_RESULT_SUCCESS) {
1602
0
            return Status::InternalError("Brotli decompression failed, result={}", result);
1603
0
        }
1604
66
        return Status::OK();
1605
66
    }
1606
};
1607
1608
Status get_block_compression_codec(segment_v2::CompressionTypePB type,
1609
34.2M
                                   BlockCompressionCodec** codec) {
1610
34.2M
    switch (type) {
1611
11.8M
    case segment_v2::CompressionTypePB::NO_COMPRESSION:
1612
11.8M
        *codec = nullptr;
1613
11.8M
        break;
1614
80.0k
    case segment_v2::CompressionTypePB::SNAPPY:
1615
80.0k
        *codec = SnappyBlockCompression::instance();
1616
80.0k
        break;
1617
98.3k
    case segment_v2::CompressionTypePB::LZ4:
1618
98.3k
        *codec = Lz4BlockCompression::instance();
1619
98.3k
        break;
1620
26.2k
    case segment_v2::CompressionTypePB::LZ4F:
1621
26.2k
        *codec = Lz4fBlockCompression::instance();
1622
26.2k
        break;
1623
3
    case segment_v2::CompressionTypePB::LZ4HC:
1624
3
        *codec = Lz4HCBlockCompression::instance();
1625
3
        break;
1626
50
    case segment_v2::CompressionTypePB::ZLIB:
1627
50
        *codec = ZlibBlockCompression::instance();
1628
50
        break;
1629
22.2M
    case segment_v2::CompressionTypePB::ZSTD:
1630
22.2M
        *codec = ZstdBlockCompression::instance();
1631
22.2M
        break;
1632
0
    default:
1633
0
        return Status::InternalError("unknown compression type({})", type);
1634
34.2M
    }
1635
1636
34.2M
    return Status::OK();
1637
34.2M
}
1638
1639
// Process-wide registry of level-aware codecs, keyed by (type, level). All
1640
// column writers that request the same codec+level share one instance, so its
1641
// internal context pool is reused according to actual write concurrency rather
1642
// than allocated once per column. Instances live for the process lifetime (like
1643
// the type-only singletons above), so their native contexts are never torn down
1644
// per segment.
1645
namespace {
1646
class LeveledCompressionCodecPool {
1647
public:
1648
271
    static LeveledCompressionCodecPool& instance() {
1649
271
        static LeveledCompressionCodecPool s_instance;
1650
271
        return s_instance;
1651
271
    }
1652
1653
270
    Status get(segment_v2::CompressionTypePB type, int level, BlockCompressionCodec** codec) {
1654
270
        const int64_t key = (static_cast<int64_t>(type) << 32) | static_cast<uint32_t>(level);
1655
270
        {
1656
270
            std::lock_guard<std::mutex> l(_mutex);
1657
270
            auto it = _codecs.find(key);
1658
270
            if (it != _codecs.end()) {
1659
261
                *codec = it->second.get();
1660
261
                return Status::OK();
1661
261
            }
1662
270
        }
1663
1664
        // Build the instance outside the lock; init() may allocate native state.
1665
9
        std::unique_ptr<BlockCompressionCodec> instance;
1666
9
        switch (type) {
1667
5
        case segment_v2::CompressionTypePB::ZSTD:
1668
5
            instance = std::make_unique<ZstdBlockCompression>(level);
1669
5
            break;
1670
4
        case segment_v2::CompressionTypePB::LZ4HC:
1671
4
            instance = std::make_unique<Lz4HCBlockCompression>(level);
1672
4
            break;
1673
0
        default:
1674
0
            return Status::InternalError("compression type({}) is not level-aware", type);
1675
9
        }
1676
9
        RETURN_IF_ERROR(instance->init());
1677
1678
9
        std::lock_guard<std::mutex> l(_mutex);
1679
        // Another thread may have inserted the same key while we were building.
1680
9
        auto [it, inserted] = _codecs.try_emplace(key, std::move(instance));
1681
9
        *codec = it->second.get();
1682
9
        return Status::OK();
1683
9
    }
1684
1685
    // Test hook: drop all pooled instances so a fresh test observes a clean pool.
1686
1
    void clear() {
1687
1
        std::lock_guard<std::mutex> l(_mutex);
1688
1
        _codecs.clear();
1689
1
    }
1690
1691
private:
1692
    std::mutex _mutex;
1693
    std::unordered_map<int64_t, std::unique_ptr<BlockCompressionCodec>> _codecs;
1694
};
1695
} // namespace
1696
1697
Status get_block_compression_codec(segment_v2::CompressionTypePB type, int level,
1698
1.05M
                                   BlockCompressionCodec** codec) {
1699
    // level <= 0 means "use codec default" -> fall back to the stateless singleton path.
1700
1.05M
    if (level <= 0) {
1701
1.05M
        return get_block_compression_codec(type, codec);
1702
1.05M
    }
1703
140
    switch (type) {
1704
137
    case segment_v2::CompressionTypePB::ZSTD:
1705
270
    case segment_v2::CompressionTypePB::LZ4HC:
1706
270
        return LeveledCompressionCodecPool::instance().get(type, level, codec);
1707
0
    default:
1708
        // types without a tunable level ignore it and use the singleton
1709
0
        return get_block_compression_codec(type, codec);
1710
140
    }
1711
140
}
1712
1713
1
void clear_leveled_compression_codec_pool_for_test() {
1714
1
    LeveledCompressionCodecPool::instance().clear();
1715
1
}
1716
1717
// this can only be used in hive text write
1718
1.10k
Status get_block_compression_codec(TFileCompressType::type type, BlockCompressionCodec** codec) {
1719
1.10k
    switch (type) {
1720
336
    case TFileCompressType::PLAIN:
1721
336
        *codec = nullptr;
1722
336
        break;
1723
0
    case TFileCompressType::ZLIB:
1724
0
        *codec = ZlibBlockCompression::instance();
1725
0
        break;
1726
153
    case TFileCompressType::GZ:
1727
153
        *codec = GzipBlockCompression::instance();
1728
153
        break;
1729
138
    case TFileCompressType::BZ2:
1730
138
        *codec = Bzip2BlockCompression::instance();
1731
138
        break;
1732
160
    case TFileCompressType::LZ4BLOCK:
1733
160
        *codec = HadoopLz4BlockCompression::instance();
1734
160
        break;
1735
158
    case TFileCompressType::SNAPPYBLOCK:
1736
158
        *codec = HadoopSnappyBlockCompression::instance();
1737
158
        break;
1738
158
    case TFileCompressType::ZSTD:
1739
158
        *codec = ZstdBlockCompression::instance();
1740
158
        break;
1741
0
    default:
1742
0
        return Status::InternalError("unsupport compression type({}) int hive text", type);
1743
1.10k
    }
1744
1745
1.10k
    return Status::OK();
1746
1.10k
}
1747
1748
Status get_block_compression_codec(tparquet::CompressionCodec::type parquet_codec,
1749
169k
                                   BlockCompressionCodec** codec) {
1750
169k
    switch (parquet_codec) {
1751
19.3k
    case tparquet::CompressionCodec::UNCOMPRESSED:
1752
19.3k
        *codec = nullptr;
1753
19.3k
        break;
1754
64.4k
    case tparquet::CompressionCodec::SNAPPY:
1755
64.4k
        *codec = SnappyBlockCompression::instance();
1756
64.4k
        break;
1757
1.78k
    case tparquet::CompressionCodec::LZ4_RAW: // we can use LZ4 compression algorithm parse LZ4_RAW
1758
2.09k
    case tparquet::CompressionCodec::LZ4:
1759
2.09k
        *codec = HadoopLz4BlockCompression::instance();
1760
2.09k
        break;
1761
73.3k
    case tparquet::CompressionCodec::ZSTD:
1762
73.3k
        *codec = ZstdBlockCompression::instance();
1763
73.3k
        break;
1764
10.0k
    case tparquet::CompressionCodec::GZIP:
1765
// Only used on x86 or x86_64
1766
10.0k
#if defined(__x86_64__) || defined(_M_X64) || defined(i386) || defined(__i386__) || \
1767
10.0k
        defined(__i386) || defined(_M_IX86)
1768
10.0k
        *codec = GzipBlockCompressionByLibdeflate::instance();
1769
#else
1770
        *codec = GzipBlockCompression::instance();
1771
#endif
1772
10.0k
        break;
1773
418
    case tparquet::CompressionCodec::LZO:
1774
418
        *codec = LzoBlockCompression::instance();
1775
418
        break;
1776
32
    case tparquet::CompressionCodec::BROTLI:
1777
32
        *codec = BrotliBlockCompression::instance();
1778
32
        break;
1779
0
    default:
1780
0
        return Status::InternalError("unknown compression type({})", parquet_codec);
1781
169k
    }
1782
1783
169k
    return Status::OK();
1784
169k
}
1785
1786
} // namespace doris