2424
2525#include " avro/DataFile.hh"
2626#include " avro/Stream.hh"
27+ #include " paimon/common/utils/options_utils.h"
2728#include " paimon/common/utils/string_utils.h"
29+ #include " paimon/core/core_options.h"
30+ #include " paimon/format/avro/avro_format_defs.h"
2831#include " paimon/format/avro/avro_format_writer.h"
2932#include " paimon/format/avro/avro_output_stream_impl.h"
3033#include " paimon/format/writer_builder.h"
@@ -58,7 +61,11 @@ class AvroWriterBuilder : public WriterBuilder {
5861 auto output_stream = std::make_unique<AvroOutputStreamImpl>(out, BUFFER_SIZE , pool_);
5962 PAIMON_ASSIGN_OR_RAISE (::avro::Codec codec,
6063 ToAvroCompressionKind (StringUtils::ToLowerCase (compression)));
61- return AvroFormatWriter::Create (std::move (output_stream), schema_, codec);
64+ PAIMON_ASSIGN_OR_RAISE (std::optional<int > compression_zstd_level,
65+ GetAvroCompressionZstdLevel (codec));
66+
67+ return AvroFormatWriter::Create (std::move (output_stream), schema_, codec,
68+ compression_zstd_level);
6269 }
6370
6471 private:
@@ -77,6 +84,14 @@ class AvroWriterBuilder : public WriterBuilder {
7784 return Status::Invalid (" unknown compression " + file_compression);
7885 }
7986 }
87+ Result<std::optional<int >> GetAvroCompressionZstdLevel (const ::avro::Codec& codec) {
88+ PAIMON_ASSIGN_OR_RAISE (CoreOptions core_options, CoreOptions::FromMap (options_));
89+ std::optional<int > compression_zstd_level;
90+ if (codec == ::avro::Codec::ZSTD_CODEC ) {
91+ compression_zstd_level = core_options.GetFileCompressionZstdLevel ();
92+ }
93+ return compression_zstd_level;
94+ }
8095
8196 std::shared_ptr<MemoryPool> pool_;
8297 std::shared_ptr<arrow::Schema> schema_;
0 commit comments