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 |