Coverage Report

Created: 2026-07-28 15:55

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