Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion be/src/storage/segment/column_writer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -486,7 +486,10 @@ ScalarColumnWriter::~ScalarColumnWriter() {
}

Status ScalarColumnWriter::init() {
RETURN_IF_ERROR(get_block_compression_codec(_opts.meta->compression(), &_compress_codec));
RETURN_IF_ERROR(get_block_compression_codec(
_opts.meta->compression(),
_opts.meta->has_compression_level() ? _opts.meta->compression_level() : 0,
&_compress_codec));

PageBuilder* page_builder = nullptr;

Expand Down
9 changes: 8 additions & 1 deletion be/src/storage/segment/segment_writer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -117,7 +117,14 @@ void SegmentWriter::init_column_meta(ColumnMetaPB* meta, uint32_t column_id,
meta->set_type(int(column.type()));
meta->set_length(column.length());
meta->set_encoding(EncodingInfo::resolve_default_encoding(opts.storage_format, column));
meta->set_compression(_opts.compression_type);
if (column.has_compression()) {
meta->set_compression(column.compression());
if (column.compression_level() > 0) {
meta->set_compression_level(column.compression_level());
}
} else {
meta->set_compression(_opts.compression_type);
}
meta->set_is_nullable(column.is_nullable());
meta->set_default_value(column.default_value());
meta->set_precision(column.precision());
Expand Down
9 changes: 8 additions & 1 deletion be/src/storage/segment/vertical_segment_writer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -134,7 +134,14 @@ void VerticalSegmentWriter::_init_column_meta(ColumnMetaPB* meta, uint32_t colum
meta->set_type(int(column.type()));
meta->set_length(cast_set<int32_t>(column.length()));
meta->set_encoding(EncodingInfo::resolve_default_encoding(opts.storage_format, column));
meta->set_compression(_opts.compression_type);
if (column.has_compression()) {
meta->set_compression(column.compression());
if (column.compression_level() > 0) {
meta->set_compression_level(column.compression_level());
}
} else {
meta->set_compression(_opts.compression_type);
}
meta->set_is_nullable(column.is_nullable());
meta->set_default_value(column.default_value());
meta->set_precision(column.precision());
Expand Down
1 change: 1 addition & 0 deletions be/src/storage/segment/vertical_segment_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,7 @@ class VerticalSegmentWriter {

private:
friend class ::doris::BlockAggregator;
friend class TestVerticalSegmentWriter;
uint32_t _segment_id;
TabletSchemaSPtr _tablet_schema;
BaseTabletSPtr _tablet;
Expand Down
28 changes: 28 additions & 0 deletions be/src/storage/tablet/tablet_meta.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -417,6 +417,34 @@ void TabletMeta::init_column_from_tcolumn(uint32_t unique_id, const TColumn& tco
if (tcolumn.__isset.variant_enable_nested_group) {
column->set_variant_enable_nested_group(tcolumn.variant_enable_nested_group);
}
if (tcolumn.__isset.compression_type) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ColumnPB 中也要设置

// The raw cast below is only valid while TCompressionType (thrift) and
// CompressionTypePB (proto) stay numerically identical. Guard every value
// at compile time so a future reorder of either enum fails to build
// instead of silently writing a wrong compression tag into segments.
static_assert(static_cast<int>(TCompressionType::UNKNOWN_COMPRESSION) ==
static_cast<int>(segment_v2::UNKNOWN_COMPRESSION));
static_assert(static_cast<int>(TCompressionType::DEFAULT_COMPRESSION) ==
static_cast<int>(segment_v2::DEFAULT_COMPRESSION));
static_assert(static_cast<int>(TCompressionType::NO_COMPRESSION) ==
static_cast<int>(segment_v2::NO_COMPRESSION));
static_assert(static_cast<int>(TCompressionType::SNAPPY) ==
static_cast<int>(segment_v2::SNAPPY));
static_assert(static_cast<int>(TCompressionType::LZ4) == static_cast<int>(segment_v2::LZ4));
static_assert(static_cast<int>(TCompressionType::LZ4F) ==
static_cast<int>(segment_v2::LZ4F));
static_assert(static_cast<int>(TCompressionType::ZLIB) ==
static_cast<int>(segment_v2::ZLIB));
static_assert(static_cast<int>(TCompressionType::ZSTD) ==
static_cast<int>(segment_v2::ZSTD));
static_assert(static_cast<int>(TCompressionType::LZ4HC) ==
static_cast<int>(segment_v2::LZ4HC));
column->set_compression_type(
static_cast<segment_v2::CompressionTypePB>(tcolumn.compression_type));
if (tcolumn.__isset.compression_level && tcolumn.compression_level > 0) {
column->set_compression_level(tcolumn.compression_level);
}
}
}

void TabletMeta::init_schema_from_thrift(const TTabletSchema& tablet_schema,
Expand Down
8 changes: 8 additions & 0 deletions be/src/storage/tablet/tablet_schema.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -555,6 +555,8 @@ void TabletColumn::init_from_pb(const ColumnPB& column) {
if (column.has_pattern_type()) {
_pattern_type = column.pattern_type();
}
_compression = column.compression_type();
_compression_level = column.has_compression_level() ? column.compression_level() : 0;
}

TabletColumn TabletColumn::create_materialized_variant_column(const std::string& root,
Expand Down Expand Up @@ -641,6 +643,12 @@ void TabletColumn::to_schema_pb(ColumnPB* column) const {
column->set_variant_doc_materialization_min_rows(_variant.doc_materialization_min_rows);
column->set_variant_doc_hash_shard_count(_variant.doc_hash_shard_count);
column->set_variant_enable_nested_group(_variant.enable_nested_group);
if (has_compression()) {
column->set_compression_type(_compression);
if (_compression_level > 0) {
column->set_compression_level(_compression_level);
}
}
}

void TabletColumn::add_sub_column(TabletColumn& sub_column) {
Expand Down
6 changes: 6 additions & 0 deletions be/src/storage/tablet/tablet_schema.h
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,9 @@ class TabletColumn : public MetadataAdder<TabletColumn> {
void set_type(FieldType type) { _type = type; }
bool is_key() const { return _is_key; }
bool is_nullable() const { return _is_nullable; }
bool has_compression() const { return _compression != segment_v2::UNKNOWN_COMPRESSION; }
segment_v2::CompressionTypePB compression() const { return _compression; }
int compression_level() const { return _compression_level; }
bool is_auto_increment() const { return _is_auto_increment; }
bool is_seqeunce_col() const { return _col_name == SEQUENCE_COL; }
bool is_on_update_current_timestamp() const { return _is_on_update_current_timestamp; }
Expand Down Expand Up @@ -302,6 +305,9 @@ class TabletColumn : public MetadataAdder<TabletColumn> {
bool _has_default_value = false;
std::string _default_value;

segment_v2::CompressionTypePB _compression = segment_v2::UNKNOWN_COMPRESSION;
int _compression_level = 0;

bool _is_decimal = false;
int32_t _precision = -1;
int32_t _frac = -1;
Expand Down
176 changes: 140 additions & 36 deletions be/src/util/block_compression.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@
#include <mutex>
#include <orc/Exceptions.hh>
#include <ostream>
#include <unordered_map>

#include "absl/strings/substitute.h"
#include "common/config.h"
Expand Down Expand Up @@ -580,6 +581,8 @@ class Lz4HCBlockCompression : public BlockCompressionCodec {
static Lz4HCBlockCompression s_instance;
return &s_instance;
}
Lz4HCBlockCompression() = default;
explicit Lz4HCBlockCompression(int level) : _compression_level(level) {}
~Lz4HCBlockCompression() {
SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
ExecEnv::GetInstance()->block_compression_mem_tracker());
Expand Down Expand Up @@ -659,10 +662,20 @@ class Lz4HCBlockCompression : public BlockCompressionCodec {
if (localCtx.get() == nullptr) {
return Status::InvalidArgument("new LZ4HC context error");
}
localCtx->ctx = LZ4_createStreamHC();
// Allocate the native stream under the compression tracker so its
// creation and the destructor's LZ4_freeStreamHC() hit the same tracker.
{
SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
ExecEnv::GetInstance()->block_compression_mem_tracker());
localCtx->ctx = LZ4_createStreamHC();
}
if (localCtx->ctx == nullptr) {
return Status::InvalidArgument("LZ4_createStreamHC error");
}
// A newly created stream defaults to the library's default level, so
// apply the requested level here; otherwise the first page compressed
// by this context would ignore the configured level.
LZ4_resetStreamHC_fast(localCtx->ctx, static_cast<int>(_compression_level));
out = std::move(localCtx);
return Status::OK();
}
Expand Down Expand Up @@ -1077,6 +1090,8 @@ class ZstdBlockCompression : public BlockCompressionCodec {
static ZstdBlockCompression s_instance;
return &s_instance;
}
ZstdBlockCompression() = default;
explicit ZstdBlockCompression(int level) : _compression_level(level) {}
~ZstdBlockCompression() {
SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
ExecEnv::GetInstance()->block_compression_mem_tracker());
Expand Down Expand Up @@ -1123,47 +1138,51 @@ class ZstdBlockCompression : public BlockCompressionCodec {
compressed_buf.size = max_len;
}

// set compression level to default 3
auto ret = ZSTD_CCtx_setParameter(context->ctx, ZSTD_c_compressionLevel,
ZSTD_CLEVEL_DEFAULT);
if (ZSTD_isError(ret)) {
return Status::InvalidArgument("ZSTD_CCtx_setParameter compression level error: {}",
ZSTD_getErrorString(ZSTD_getErrorCode(ret)));
}
// set checksum flag to 1
ret = ZSTD_CCtx_setParameter(context->ctx, ZSTD_c_checksumFlag, 1);
if (ZSTD_isError(ret)) {
return Status::InvalidArgument("ZSTD_CCtx_setParameter checksumFlag error: {}",
ZSTD_getErrorString(ZSTD_getErrorCode(ret)));
}

ZSTD_outBuffer out_buf = {compressed_buf.data, compressed_buf.size, 0};
{
SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
ExecEnv::GetInstance()->block_compression_mem_tracker());
auto ret = ZSTD_CCtx_setParameter(context->ctx, ZSTD_c_compressionLevel,
_compression_level);
if (ZSTD_isError(ret)) {
return Status::InvalidArgument(
"ZSTD_CCtx_setParameter compression level error: {}",
ZSTD_getErrorString(ZSTD_getErrorCode(ret)));
}
// set checksum flag to 1
ret = ZSTD_CCtx_setParameter(context->ctx, ZSTD_c_checksumFlag, 1);
if (ZSTD_isError(ret)) {
return Status::InvalidArgument("ZSTD_CCtx_setParameter checksumFlag error: {}",
ZSTD_getErrorString(ZSTD_getErrorCode(ret)));
}

for (size_t i = 0; i < inputs.size(); i++) {
ZSTD_inBuffer in_buf = {inputs[i].data, inputs[i].size, 0};
for (size_t i = 0; i < inputs.size(); i++) {
ZSTD_inBuffer in_buf = {inputs[i].data, inputs[i].size, 0};

bool last_input = (i == inputs.size() - 1);
auto mode = last_input ? ZSTD_e_end : ZSTD_e_continue;
bool last_input = (i == inputs.size() - 1);
auto mode = last_input ? ZSTD_e_end : ZSTD_e_continue;

bool finished = false;
do {
// do compress
ret = ZSTD_compressStream2(context->ctx, &out_buf, &in_buf, mode);
bool finished = false;
do {
// do compress
ret = ZSTD_compressStream2(context->ctx, &out_buf, &in_buf, mode);

if (ZSTD_isError(ret)) {
compress_failed = true;
return Status::InternalError("ZSTD_compressStream2 error: {}",
ZSTD_getErrorString(ZSTD_getErrorCode(ret)));
}
if (ZSTD_isError(ret)) {
compress_failed = true;
return Status::InternalError(
"ZSTD_compressStream2 error: {}",
ZSTD_getErrorString(ZSTD_getErrorCode(ret)));
}

// ret is ZSTD hint for needed output buffer size
if (ret > 0 && out_buf.pos == out_buf.size) {
compress_failed = true;
return Status::InternalError("ZSTD_compressStream2 output buffer full");
}
// ret is ZSTD hint for needed output buffer size
if (ret > 0 && out_buf.pos == out_buf.size) {
compress_failed = true;
return Status::InternalError("ZSTD_compressStream2 output buffer full");
}

finished = last_input ? (ret == 0) : (in_buf.pos == inputs[i].size);
} while (!finished);
finished = last_input ? (ret == 0) : (in_buf.pos == inputs[i].size);
} while (!finished);
}
}

// set compressed size for caller
Expand Down Expand Up @@ -1215,7 +1234,13 @@ class ZstdBlockCompression : public BlockCompressionCodec {
return Status::InvalidArgument("failed to new ZSTD CContext");
}
//typedef LZ4F_cctx* LZ4F_compressionContext_t;
localCtx->ctx = ZSTD_createCCtx();
// Allocate the native context under the compression tracker so its
// creation and the destructor's ZSTD_freeCCtx() hit the same tracker.
{
SCOPED_SWITCH_THREAD_MEM_TRACKER_LIMITER(
ExecEnv::GetInstance()->block_compression_mem_tracker());
localCtx->ctx = ZSTD_createCCtx();
}
if (localCtx->ctx == nullptr) {
return Status::InvalidArgument("Failed to create ZSTD compress ctx");
}
Expand Down Expand Up @@ -1262,6 +1287,7 @@ class ZstdBlockCompression : public BlockCompressionCodec {
}

private:
int _compression_level = ZSTD_CLEVEL_DEFAULT;
mutable std::mutex _ctx_c_mutex;
mutable std::vector<std::unique_ptr<CContext>> _ctx_c_pool;

Expand Down Expand Up @@ -1615,6 +1641,84 @@ Status get_block_compression_codec(segment_v2::CompressionTypePB type,
return Status::OK();
}

// Process-wide registry of level-aware codecs, keyed by (type, level). All
// column writers that request the same codec+level share one instance, so its
// internal context pool is reused according to actual write concurrency rather
// than allocated once per column. Instances live for the process lifetime (like
// the type-only singletons above), so their native contexts are never torn down
// per segment.
namespace {
class LeveledCompressionCodecPool {
public:
static LeveledCompressionCodecPool& instance() {
static LeveledCompressionCodecPool s_instance;
return s_instance;
}

Status get(segment_v2::CompressionTypePB type, int level, BlockCompressionCodec** codec) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

这里如果是旧表的话会传入 level = 0,之前旧表的默认的 ZSTD 或者 LZ4HC 的 level 是 0 吗

const int64_t key = (static_cast<int64_t>(type) << 32) | static_cast<uint32_t>(level);
{
std::lock_guard<std::mutex> l(_mutex);
auto it = _codecs.find(key);
if (it != _codecs.end()) {
*codec = it->second.get();
return Status::OK();
}
}

// Build the instance outside the lock; init() may allocate native state.
std::unique_ptr<BlockCompressionCodec> instance;
switch (type) {
case segment_v2::CompressionTypePB::ZSTD:
instance = std::make_unique<ZstdBlockCompression>(level);
break;
case segment_v2::CompressionTypePB::LZ4HC:
instance = std::make_unique<Lz4HCBlockCompression>(level);
break;
default:
return Status::InternalError("compression type({}) is not level-aware", type);
}
RETURN_IF_ERROR(instance->init());

std::lock_guard<std::mutex> l(_mutex);
// Another thread may have inserted the same key while we were building.
auto it = _codecs.try_emplace(key, std::move(instance)).first;
*codec = it->second.get();
return Status::OK();
}

// Test hook: drop all pooled instances so a fresh test observes a clean pool.
void clear() {
std::lock_guard<std::mutex> l(_mutex);
_codecs.clear();
}

private:
std::mutex _mutex;
std::unordered_map<int64_t, std::unique_ptr<BlockCompressionCodec>> _codecs;
};
} // namespace

Status get_block_compression_codec(segment_v2::CompressionTypePB type, int level,
BlockCompressionCodec** codec) {
// level <= 0 means "use codec default" -> fall back to the stateless singleton path.
if (level <= 0) {
return get_block_compression_codec(type, codec);
}
switch (type) {
case segment_v2::CompressionTypePB::ZSTD:
case segment_v2::CompressionTypePB::LZ4HC:
return LeveledCompressionCodecPool::instance().get(type, level, codec);
default:
// types without a tunable level ignore it and use the singleton
return get_block_compression_codec(type, codec);
}
}

void clear_leveled_compression_codec_pool_for_test() {
LeveledCompressionCodecPool::instance().clear();
}

// this can only be used in hive text write
Status get_block_compression_codec(TFileCompressType::type type, BlockCompressionCodec** codec) {
switch (type) {
Expand Down
10 changes: 10 additions & 0 deletions be/src/util/block_compression.h
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,16 @@ class BlockCompressionCodec {
Status get_block_compression_codec(segment_v2::CompressionTypePB type,
BlockCompressionCodec** codec);

// Level-aware variant. If `level` > 0 and `type` is a level-aware codec (ZSTD, LZ4HC),
// returns a process-wide instance shared by all callers that request the same
// (type, level) pair (do not delete). Otherwise `*codec` points to the type-only
// singleton. In both cases the returned codec is owned by the process, not the caller.
Status get_block_compression_codec(segment_v2::CompressionTypePB type, int level,
BlockCompressionCodec** codec);

// Test-only: drops all pooled level-aware codec instances.
void clear_leveled_compression_codec_pool_for_test();

Status get_block_compression_codec(tparquet::CompressionCodec::type parquet_codec,
BlockCompressionCodec** codec);

Expand Down
Loading
Loading