Skip to content

Commit 61d6560

Browse files
committed
feat(avro): support zstd compression level configuration
1 parent 623db38 commit 61d6560

8 files changed

Lines changed: 108 additions & 11 deletions

src/paimon/format/avro/avro_file_batch_reader_test.cpp

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,8 @@ namespace paimon::avro::test {
3939
class AvroFileBatchReaderTest : public ::testing::Test, public ::testing::WithParamInterface<bool> {
4040
public:
4141
void SetUp() override {
42-
ASSERT_OK_AND_ASSIGN(file_format_, FileFormatFactory::Get("avro", {}));
42+
ASSERT_OK_AND_ASSIGN(file_format_,
43+
FileFormatFactory::Get("avro", {{Options::FILE_FORMAT, "avro"}}));
4344
fs_ = std::make_shared<LocalFileSystem>();
4445
dir_ = ::paimon::test::UniqueTestDirectory::Create();
4546
ASSERT_TRUE(dir_);

src/paimon/format/avro/avro_file_format_test.cpp

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,8 @@ namespace paimon::avro::test {
4545
class AvroFileFormatTest : public testing::Test, public ::testing::WithParamInterface<std::string> {
4646
public:
4747
void SetUp() override {
48-
ASSERT_OK_AND_ASSIGN(file_format_, FileFormatFactory::Get("avro", {}));
48+
ASSERT_OK_AND_ASSIGN(file_format_,
49+
FileFormatFactory::Get("avro", {{Options::FILE_FORMAT, "avro"}}));
4950
fs_ = std::make_shared<LocalFileSystem>();
5051
dir_ = ::paimon::test::UniqueTestDirectory::Create();
5152
ASSERT_TRUE(dir_);

src/paimon/format/avro/avro_format_writer.cpp

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -53,13 +53,14 @@ AvroFormatWriter::AvroFormatWriter(std::unique_ptr<::avro::DataFileWriterBase>&&
5353

5454
Result<std::unique_ptr<AvroFormatWriter>> AvroFormatWriter::Create(
5555
std::unique_ptr<AvroOutputStreamImpl> out, const std::shared_ptr<arrow::Schema>& schema,
56-
const ::avro::Codec codec) {
56+
const ::avro::Codec codec, std::optional<int> compression_level) {
5757
try {
5858
PAIMON_ASSIGN_OR_RAISE(::avro::ValidSchema avro_schema,
5959
AvroSchemaConverter::ArrowSchemaToAvroSchema(schema));
6060
AvroOutputStreamImpl* avro_output_stream = out.get();
61-
auto writer = std::make_unique<::avro::DataFileWriterBase>(std::move(out), avro_schema,
62-
DEFAULT_SYNC_INTERVAL, codec);
61+
auto writer = std::make_unique<::avro::DataFileWriterBase>(
62+
std::move(out), avro_schema, DEFAULT_SYNC_INTERVAL, codec, ::avro::Metadata(),
63+
compression_level);
6364
auto data_type = arrow::struct_(schema->fields());
6465
return std::unique_ptr<AvroFormatWriter>(
6566
new AvroFormatWriter(std::move(writer), avro_schema, data_type, avro_output_stream));

src/paimon/format/avro/avro_format_writer.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
#include <cstddef>
2020
#include <cstdint>
2121
#include <memory>
22+
#include <optional>
2223

2324
#include "arrow/api.h"
2425
#include "avro/DataFile.hh"
@@ -49,7 +50,7 @@ class AvroFormatWriter : public FormatWriter {
4950
public:
5051
static Result<std::unique_ptr<AvroFormatWriter>> Create(
5152
std::unique_ptr<AvroOutputStreamImpl> out, const std::shared_ptr<arrow::Schema>& schema,
52-
const ::avro::Codec codec);
53+
const ::avro::Codec codec, std::optional<int> compression_level = std::nullopt);
5354

5455
Status AddBatch(ArrowArray* batch) override;
5556

src/paimon/format/avro/avro_format_writer_test.cpp

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -68,7 +68,8 @@ class AvroFormatWriterTest : public ::testing::Test {
6868
int32_t batch_size) {
6969
::ArrowSchema c_schema;
7070
EXPECT_TRUE(arrow::ExportSchema(*schema, &c_schema).ok());
71-
EXPECT_OK_AND_ASSIGN(auto file_format, FileFormatFactory::Get("avro", {}));
71+
EXPECT_OK_AND_ASSIGN(auto file_format,
72+
FileFormatFactory::Get("avro", {{Options::FILE_FORMAT, "avro"}}));
7273
EXPECT_OK_AND_ASSIGN(auto writer_builder,
7374
file_format->CreateWriterBuilder(&c_schema, batch_size));
7475
EXPECT_OK_AND_ASSIGN(std::shared_ptr<FormatWriter> writer,

src/paimon/format/avro/avro_writer_builder.h

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
#include "avro/DataFile.hh"
2626
#include "avro/Stream.hh"
2727
#include "paimon/common/utils/string_utils.h"
28+
#include "paimon/core/core_options.h"
2829
#include "paimon/format/avro/avro_format_writer.h"
2930
#include "paimon/format/avro/avro_output_stream_impl.h"
3031
#include "paimon/format/writer_builder.h"
@@ -56,13 +57,19 @@ class AvroWriterBuilder : public WriterBuilder {
5657
Result<std::unique_ptr<FormatWriter>> Build(const std::shared_ptr<OutputStream>& out,
5758
const std::string& compression) override {
5859
auto output_stream = std::make_unique<AvroOutputStreamImpl>(out, BUFFER_SIZE, pool_);
60+
std::string file_compression =
61+
options_.find(AVRO_CODEC) != options_.end() ? options_.at(AVRO_CODEC) : compression;
5962
PAIMON_ASSIGN_OR_RAISE(::avro::Codec codec,
60-
ToAvroCompressionKind(StringUtils::ToLowerCase(compression)));
61-
return AvroFormatWriter::Create(std::move(output_stream), schema_, codec);
63+
ToAvroCompressionKind(StringUtils::ToLowerCase(file_compression)));
64+
PAIMON_ASSIGN_OR_RAISE(std::optional<int> compression_level,
65+
GetAvroCompressionLevel(codec));
66+
return AvroFormatWriter::Create(std::move(output_stream), schema_, codec,
67+
compression_level);
6268
}
6369

6470
private:
6571
static constexpr int32_t BUFFER_SIZE = 1024 * 1024;
72+
static inline const char AVRO_CODEC[] = "avro.codec";
6673

6774
static Result<::avro::Codec> ToAvroCompressionKind(const std::string& file_compression) {
6875
if (file_compression == "zstd" || file_compression == "zstandard") {
@@ -77,6 +84,14 @@ class AvroWriterBuilder : public WriterBuilder {
7784
return Status::Invalid("unknown compression " + file_compression);
7885
}
7986
}
87+
Result<std::optional<int>> GetAvroCompressionLevel(const ::avro::Codec& codec) {
88+
std::optional<int> compression_level;
89+
if (codec == ::avro::Codec::ZSTD_CODEC) {
90+
PAIMON_ASSIGN_OR_RAISE(CoreOptions core_options, CoreOptions::FromMap(options_));
91+
compression_level = core_options.GetFileCompressionZstdLevel();
92+
}
93+
return compression_level;
94+
}
8095

8196
std::shared_ptr<MemoryPool> pool_;
8297
std::shared_ptr<arrow::Schema> schema_;

src/paimon/format/avro/avro_writer_builder_test.cpp

Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,4 +51,81 @@ TEST(ToAvroCompressionKindTest, HandlesInvalidCompression) {
5151
TEST(ToAvroCompressionKindTest, HandlesEmptyString) {
5252
ASSERT_NOK(AvroWriterBuilder::ToAvroCompressionKind(""));
5353
}
54+
55+
TEST(ToAvroCompressionKindTest, CheckAvroCodec) {
56+
arrow::FieldVector fields = {arrow::field("f0", arrow::int32())};
57+
auto schema = std::make_shared<arrow::Schema>(fields);
58+
AvroWriterBuilder builder(schema, -1,
59+
{{Options::FILE_FORMAT, "avro"}, {"avro.codec", "snappy"}});
60+
ASSERT_OK_AND_ASSIGN(auto file_writer, builder.Build(nullptr, "zstd"));
61+
auto* avro_file_writer = dynamic_cast<AvroFormatWriter*>(file_writer.get());
62+
ASSERT_EQ(avro_file_writer->writer_->codec_, ::avro::Codec::SNAPPY_CODEC);
63+
ASSERT_EQ(avro_file_writer->writer_->compressionLevel_, std::nullopt);
64+
65+
AvroWriterBuilder builder2(schema, -1,
66+
{{Options::FILE_FORMAT, "avro"}, {"avro.codec", "deflate"}});
67+
ASSERT_OK_AND_ASSIGN(auto file_writer2, builder2.Build(nullptr, "zstd"));
68+
auto* avro_file_writer2 = dynamic_cast<AvroFormatWriter*>(file_writer2.get());
69+
ASSERT_EQ(avro_file_writer2->writer_->codec_, ::avro::Codec::DEFLATE_CODEC);
70+
ASSERT_EQ(avro_file_writer2->writer_->compressionLevel_, std::nullopt);
71+
72+
AvroWriterBuilder builder3(schema, -1,
73+
{{Options::FILE_FORMAT, "avro"}, {"avro.codec", "zstd"}});
74+
ASSERT_OK_AND_ASSIGN(auto file_writer3, builder3.Build(nullptr, "zstd"));
75+
auto* avro_file_writer3 = dynamic_cast<AvroFormatWriter*>(file_writer3.get());
76+
ASSERT_EQ(avro_file_writer3->writer_->codec_, ::avro::Codec::ZSTD_CODEC);
77+
ASSERT_EQ(avro_file_writer3->writer_->compressionLevel_, 1);
78+
79+
AvroWriterBuilder builder4(schema, -1,
80+
{{Options::FILE_FORMAT, "avro"},
81+
{"avro.codec", "zstd"},
82+
{Options::FILE_COMPRESSION_ZSTD_LEVEL, "3"}});
83+
ASSERT_OK_AND_ASSIGN(auto file_writer4, builder4.Build(nullptr, "zstd"));
84+
auto* avro_file_writer4 = dynamic_cast<AvroFormatWriter*>(file_writer4.get());
85+
ASSERT_EQ(avro_file_writer4->writer_->codec_, ::avro::Codec::ZSTD_CODEC);
86+
ASSERT_EQ(avro_file_writer4->writer_->compressionLevel_, 3);
87+
88+
AvroWriterBuilder builder5(schema, -1,
89+
{{Options::FILE_FORMAT, "avro"},
90+
{"avro.codec", "null"},
91+
{Options::FILE_COMPRESSION_ZSTD_LEVEL, "3"}});
92+
ASSERT_OK_AND_ASSIGN(auto file_writer5, builder5.Build(nullptr, "zstd"));
93+
auto* avro_file_writer5 = dynamic_cast<AvroFormatWriter*>(file_writer5.get());
94+
ASSERT_EQ(avro_file_writer5->writer_->codec_, ::avro::Codec::NULL_CODEC);
95+
ASSERT_EQ(avro_file_writer5->writer_->compressionLevel_, std::nullopt);
96+
97+
AvroWriterBuilder builder6(schema, -1,
98+
{{Options::FILE_FORMAT, "avro"},
99+
{"avro.codec", "test"},
100+
{Options::FILE_COMPRESSION_ZSTD_LEVEL, "3"}});
101+
ASSERT_NOK(builder6.Build(nullptr, "zstd"));
102+
}
103+
104+
TEST(ToAvroCompressionKindTest, CheckAvroCompressionLevel) {
105+
AvroWriterBuilder builder(nullptr, -1, {{Options::FILE_FORMAT, "avro"}});
106+
ASSERT_OK_AND_ASSIGN(std::optional<int> zstd_level,
107+
builder.GetAvroCompressionLevel(::avro::Codec::ZSTD_CODEC));
108+
ASSERT_TRUE(zstd_level.has_value());
109+
ASSERT_EQ(zstd_level.value(), 1);
110+
111+
ASSERT_OK_AND_ASSIGN(std::optional<int> compression_level1,
112+
builder.GetAvroCompressionLevel(::avro::Codec::SNAPPY_CODEC));
113+
ASSERT_FALSE(compression_level1.has_value());
114+
115+
ASSERT_OK_AND_ASSIGN(std::optional<int> compression_level2,
116+
builder.GetAvroCompressionLevel(::avro::Codec::DEFLATE_CODEC));
117+
ASSERT_FALSE(compression_level2.has_value());
118+
119+
ASSERT_OK_AND_ASSIGN(std::optional<int> compression_level3,
120+
builder.GetAvroCompressionLevel(::avro::Codec::NULL_CODEC));
121+
ASSERT_FALSE(compression_level3.has_value());
122+
123+
AvroWriterBuilder builder2(
124+
nullptr, -1, {{Options::FILE_FORMAT, "avro"}, {Options::FILE_COMPRESSION_ZSTD_LEVEL, "3"}});
125+
ASSERT_OK_AND_ASSIGN(std::optional<int> zstd_level2,
126+
builder2.GetAvroCompressionLevel(::avro::Codec::ZSTD_CODEC));
127+
ASSERT_TRUE(zstd_level2.has_value());
128+
ASSERT_EQ(zstd_level2.value(), 3);
129+
}
130+
54131
} // namespace paimon::avro::test

third_party/versions.txt

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -55,8 +55,8 @@ PAIMON_GTEST_PKG_NAME=gtest-${PAIMON_GTEST_BUILD_VERSION}.tar.gz
5555
PAIMON_ARROW_BUILD_VERSION=17.0.0
5656
PAIMON_ARROW_BUILD_SHA256_CHECKSUM=9d280d8042e7cf526f8c28d170d93bfab65e50f94569f6a790982a878d8d898d
5757
PAIMON_ARROW_PKG_NAME=apache-arrow-${PAIMON_ARROW_BUILD_VERSION}.tar.gz
58-
PAIMON_AVRO_BUILD_VERSION=54b332161524086dcb6cde8afe097097eed7f3ee
59-
PAIMON_AVRO_BUILD_SHA256_CHECKSUM=00febd590b1e328d3a97b67a6d29a1d0243e0e41bb2b1582ec580d37698d1fe2
58+
PAIMON_AVRO_BUILD_VERSION=c499eefb48aa2db906c7bca14a047223806f36db
59+
PAIMON_AVRO_BUILD_SHA256_CHECKSUM=9771f1dcfe3c01aff7ff670e873e66d3406362f71941821d482de65f3d32d780
6060
PAIMON_AVRO_PKG_NAME=avro-${PAIMON_AVRO_BUILD_VERSION}.tar.gz
6161
PAIMON_FMT_BUILD_VERSION=11.2.0
6262
PAIMON_FMT_BUILD_SHA256_CHECKSUM=bc23066d87ab3168f27cef3e97d545fa63314f5c79df5ea444d41d56f962c6af

0 commit comments

Comments
 (0)