be/src/format/transformer/vparquet_writer.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 "format/transformer/vparquet_writer.h" |
19 | | |
20 | | #include <arrow/io/type_fwd.h> |
21 | | #include <arrow/table.h> |
22 | | #include <glog/logging.h> |
23 | | #include <parquet/api/reader.h> |
24 | | #include <parquet/column_writer.h> |
25 | | #include <parquet/platform.h> |
26 | | #include <parquet/schema.h> |
27 | | #include <parquet/type_fwd.h> |
28 | | #include <parquet/types.h> |
29 | | |
30 | | #include <ctime> |
31 | | #include <exception> |
32 | | #include <ostream> |
33 | | #include <string> |
34 | | |
35 | | #include "common/config.h" |
36 | | #include "common/status.h" |
37 | | #include "exprs/vexpr.h" |
38 | | #include "exprs/vexpr_context.h" |
39 | | #include "format/arrow/arrow_row_batch.h" |
40 | | #include "format/arrow/arrow_utils.h" |
41 | | #include "format/parquet/parquet_arrow_block_convertor.h" |
42 | | #include "io/fs/file_writer.h" |
43 | | #include "runtime/exec_env.h" |
44 | | #include "runtime/runtime_state.h" |
45 | | #include "util/debug_util.h" |
46 | | #include "util/timezone_utils.h" |
47 | | |
48 | | namespace doris { |
49 | | |
50 | | ParquetOutputStream::ParquetOutputStream(doris::io::FileWriter* file_writer) |
51 | 19 | : _file_writer(file_writer), _cur_pos(0), _written_len(0) { |
52 | 19 | set_mode(arrow::io::FileMode::WRITE); |
53 | 19 | } Unexecuted instantiation: _ZN5doris19ParquetOutputStreamC2EPNS_2io10FileWriterE _ZN5doris19ParquetOutputStreamC1EPNS_2io10FileWriterE Line | Count | Source | 51 | 19 | : _file_writer(file_writer), _cur_pos(0), _written_len(0) { | 52 | 19 | set_mode(arrow::io::FileMode::WRITE); | 53 | 19 | } |
|
54 | | |
55 | 19 | ParquetOutputStream::~ParquetOutputStream() { |
56 | 19 | arrow::Status st = Close(); |
57 | 19 | if (!st.ok()) { |
58 | 0 | LOG(WARNING) << "close parquet file error: " << st.ToString(); |
59 | 0 | } |
60 | 19 | } |
61 | | |
62 | 1.20k | arrow::Status ParquetOutputStream::Write(const void* data, int64_t nbytes) { |
63 | 1.20k | if (_is_closed) { |
64 | 0 | return arrow::Status::OK(); |
65 | 0 | } |
66 | 1.20k | size_t written_len = nbytes; |
67 | 1.20k | Status st = _file_writer->append({static_cast<const uint8_t*>(data), written_len}); |
68 | 1.20k | if (!st.ok()) { |
69 | 0 | return arrow::Status::IOError(st.to_string()); |
70 | 0 | } |
71 | 1.20k | _cur_pos += written_len; |
72 | 1.20k | _written_len += written_len; |
73 | 1.20k | return arrow::Status::OK(); |
74 | 1.20k | } |
75 | | |
76 | 1.94k | arrow::Result<int64_t> ParquetOutputStream::Tell() const { |
77 | 1.94k | return _cur_pos; |
78 | 1.94k | } |
79 | | |
80 | 38 | arrow::Status ParquetOutputStream::Close() { |
81 | 38 | if (!_is_closed) { |
82 | 19 | Defer defer {[this] { _is_closed = true; }}; |
83 | 19 | Status st = _file_writer->close(); |
84 | 19 | if (!st.ok()) { |
85 | 0 | LOG(WARNING) << "close parquet output stream failed: " << st; |
86 | 0 | return arrow::Status::IOError(st.to_string()); |
87 | 0 | } |
88 | 19 | } |
89 | 38 | return arrow::Status::OK(); |
90 | 38 | } |
91 | | |
92 | 5 | int64_t ParquetOutputStream::get_written_len() const { |
93 | 5 | return _written_len; |
94 | 5 | } |
95 | | |
96 | 0 | void ParquetOutputStream::set_written_len(int64_t written_len) { |
97 | 0 | _written_len = written_len; |
98 | 0 | } |
99 | | |
100 | | void ParquetBuildHelper::build_compression_type( |
101 | | ::parquet::WriterProperties::Builder& builder, |
102 | 19 | const TParquetCompressionType::type& compression_type) { |
103 | 19 | switch (compression_type) { |
104 | 0 | case TParquetCompressionType::SNAPPY: { |
105 | 0 | builder.compression(arrow::Compression::SNAPPY); |
106 | 0 | break; |
107 | 0 | } |
108 | 0 | case TParquetCompressionType::GZIP: { |
109 | 0 | builder.compression(arrow::Compression::GZIP); |
110 | 0 | break; |
111 | 0 | } |
112 | 0 | case TParquetCompressionType::BROTLI: { |
113 | 0 | builder.compression(arrow::Compression::BROTLI); |
114 | 0 | break; |
115 | 0 | } |
116 | 0 | case TParquetCompressionType::ZSTD: { |
117 | 0 | builder.compression(arrow::Compression::ZSTD); |
118 | 0 | break; |
119 | 0 | } |
120 | 0 | case TParquetCompressionType::LZ4: { |
121 | 0 | builder.compression(arrow::Compression::LZ4); |
122 | 0 | break; |
123 | 0 | } |
124 | 0 | case TParquetCompressionType::LZ4_HADOOP: { |
125 | 0 | constexpr int64_t HADOOP_LZ4_DEFAULT_BUFFER_SIZE = 256 * 1024; |
126 | | // Hadoop-framed LZ4 -> Parquet thrift codec "LZ4" (deprecated). This matches what |
127 | | // Spark/Iceberg writes for `write.parquet.compression-codec=lz4`. Arrow 17 emits one |
128 | | // Hadoop LZ4 block per Parquet page/dictionary page, while Hadoop JVM readers default to |
129 | | // a 256 KiB LZ4 codec buffer, so keep page targets below that buffer size. |
130 | 0 | builder.compression(arrow::Compression::LZ4_HADOOP); |
131 | 0 | builder.data_pagesize(HADOOP_LZ4_DEFAULT_BUFFER_SIZE / 2); |
132 | 0 | builder.dictionary_pagesize_limit(HADOOP_LZ4_DEFAULT_BUFFER_SIZE / 2); |
133 | 0 | break; |
134 | 0 | } |
135 | | // arrow do not support lzo and bz2 compression type. |
136 | | // case TParquetCompressionType::LZO: { |
137 | | // builder.compression(arrow::Compression::LZO); |
138 | | // break; |
139 | | // } |
140 | | // case TParquetCompressionType::BZ2: { |
141 | | // builder.compression(arrow::Compression::BZ2); |
142 | | // break; |
143 | | // } |
144 | 19 | case TParquetCompressionType::UNCOMPRESSED: { |
145 | 19 | builder.compression(arrow::Compression::UNCOMPRESSED); |
146 | 19 | break; |
147 | 0 | } |
148 | 0 | default: |
149 | 0 | builder.compression(arrow::Compression::SNAPPY); |
150 | 19 | } |
151 | 19 | } |
152 | | |
153 | | void ParquetBuildHelper::build_version(::parquet::WriterProperties::Builder& builder, |
154 | 19 | const TParquetVersion::type& parquet_version) { |
155 | 19 | switch (parquet_version) { |
156 | 18 | case TParquetVersion::PARQUET_1_0: { |
157 | 18 | builder.version(::parquet::ParquetVersion::PARQUET_1_0); |
158 | 18 | break; |
159 | 0 | } |
160 | 1 | case TParquetVersion::PARQUET_2_LATEST: { |
161 | 1 | builder.version(::parquet::ParquetVersion::PARQUET_2_LATEST); |
162 | 1 | break; |
163 | 0 | } |
164 | 0 | default: |
165 | 0 | builder.version(::parquet::ParquetVersion::PARQUET_1_0); |
166 | 19 | } |
167 | 19 | } |
168 | | |
169 | | VParquetWriter::VParquetWriter(RuntimeState* state, doris::io::FileWriter* file_writer, |
170 | | const VExprContextSPtrs& output_vexpr_ctxs, |
171 | | std::vector<std::string> column_names, bool output_object_data, |
172 | | const ParquetFileOptions& parquet_options) |
173 | 19 | : VFileFormatTransformer(state, output_vexpr_ctxs, output_object_data), |
174 | 19 | _column_names(std::move(column_names)), |
175 | 19 | _parquet_options(parquet_options) { |
176 | 19 | _outstream = std::shared_ptr<ParquetOutputStream>(new ParquetOutputStream(file_writer)); |
177 | 19 | } |
178 | | |
179 | | VParquetWriter::VParquetWriter(RuntimeState* state, doris::io::FileWriter* file_writer, |
180 | | const VExprContextSPtrs& output_vexpr_ctxs, |
181 | | std::vector<TParquetSchema> parquet_schemas, bool output_object_data, |
182 | | const ParquetFileOptions& parquet_options) |
183 | 0 | : VFileFormatTransformer(state, output_vexpr_ctxs, output_object_data), |
184 | 0 | _parquet_schemas(std::move(parquet_schemas)), |
185 | 0 | _parquet_options(parquet_options) { |
186 | 0 | _outstream = std::shared_ptr<ParquetOutputStream>(new ParquetOutputStream(file_writer)); |
187 | 0 | } |
188 | | |
189 | 19 | Status VParquetWriter::_parse_properties() { |
190 | 19 | try { |
191 | 19 | arrow::MemoryPool* pool = ExecEnv::GetInstance()->arrow_memory_pool(); |
192 | | |
193 | | //build parquet writer properties |
194 | 19 | ::parquet::WriterProperties::Builder builder; |
195 | 19 | ParquetBuildHelper::build_compression_type(builder, _parquet_options.compression_type); |
196 | 19 | ParquetBuildHelper::build_version(builder, _parquet_options.parquet_version); |
197 | 19 | if (_parquet_options.parquet_disable_dictionary) { |
198 | 0 | builder.disable_dictionary(); |
199 | 19 | } else { |
200 | 19 | builder.enable_dictionary(); |
201 | 19 | } |
202 | 19 | builder.created_by( |
203 | 19 | fmt::format("{}({})", doris::get_short_version(), ::parquet::DEFAULT_CREATED_BY)); |
204 | 19 | builder.max_row_group_length(std::numeric_limits<int64_t>::max()); |
205 | 19 | builder.memory_pool(pool); |
206 | 19 | _parquet_writer_properties = builder.build(); |
207 | | |
208 | | //build arrow writer properties |
209 | 19 | ::parquet::ArrowWriterProperties::Builder arrow_builder; |
210 | 19 | if (_parquet_options.enable_int96_timestamps) { |
211 | 6 | arrow_builder.enable_force_write_int96_timestamps(); |
212 | 6 | } |
213 | 19 | arrow_builder.store_schema(); |
214 | 19 | _arrow_properties = arrow_builder.build(); |
215 | 19 | } catch (const ::parquet::ParquetException& e) { |
216 | 0 | return Status::InternalError("parquet writer parse properties error: {}", e.what()); |
217 | 0 | } |
218 | 19 | return Status::OK(); |
219 | 19 | } |
220 | | |
221 | | std::unique_ptr<ArrowBlockConvertor> VParquetWriter::_create_arrow_block_convertor( |
222 | | DataTypes types, std::vector<std::string> names, const std::string& timezone_name, |
223 | 2 | const cctz::time_zone& timezone, bool enable_int96_timestamps) const { |
224 | 2 | return std::make_unique<ParquetArrowBlockConvertor>( |
225 | 2 | std::move(types), std::move(names), timezone_name, timezone, enable_int96_timestamps); |
226 | 2 | } |
227 | | |
228 | 19 | Status VParquetWriter::write(const Block& block) { |
229 | 19 | if (block.rows() == 0) { |
230 | 0 | return Status::OK(); |
231 | 0 | } |
232 | | |
233 | | // serialize |
234 | 19 | std::shared_ptr<arrow::RecordBatch> result; |
235 | 19 | RETURN_IF_ERROR(_arrow_block_convertor->convert_to_arrow( |
236 | 19 | block, ExecEnv::GetInstance()->arrow_memory_pool(), &result)); |
237 | 19 | if (_write_size == 0) { |
238 | 19 | RETURN_DORIS_STATUS_IF_ERROR(_writer->NewBufferedRowGroup()); |
239 | 19 | } |
240 | 19 | RETURN_DORIS_STATUS_IF_ERROR(_writer->WriteRecordBatch(*result)); |
241 | 19 | _write_size += block.bytes(); |
242 | 19 | if (_write_size >= doris::config::min_row_group_size) { |
243 | 0 | _write_size = 0; |
244 | 0 | } |
245 | 19 | return Status::OK(); |
246 | 19 | } |
247 | | |
248 | 19 | arrow::Status VParquetWriter::_open_file_writer() { |
249 | 19 | ARROW_ASSIGN_OR_RAISE(_writer, ::parquet::arrow::FileWriter::Open( |
250 | 19 | *_arrow_block_convertor->arrow_schema(), |
251 | 19 | ExecEnv::GetInstance()->arrow_memory_pool(), _outstream, |
252 | 19 | _parquet_writer_properties, _arrow_properties)); |
253 | 19 | return arrow::Status::OK(); |
254 | 19 | } |
255 | | |
256 | 19 | Status VParquetWriter::open() { |
257 | 19 | _timezone = _state->timezone(); |
258 | 19 | _timezone_obj = _state->timezone_obj(); |
259 | 19 | if (_parquet_options.enable_int96_timestamps && _parquet_options.int96_timezone.has_value()) { |
260 | 4 | _timezone = *_parquet_options.int96_timezone; |
261 | | // Cache the override on this writer, never mutate the shared query RuntimeState. |
262 | 4 | if (!TimezoneUtils::find_cctz_time_zone(_timezone, _timezone_obj)) { |
263 | 0 | return Status::InvalidArgument("Invalid Parquet INT96 writer timezone: {}", _timezone); |
264 | 0 | } |
265 | 4 | } |
266 | 19 | RETURN_IF_ERROR(_parse_properties()); |
267 | 19 | DataTypes types; |
268 | 19 | types.reserve(_output_vexpr_ctxs.size()); |
269 | 378 | for (const auto& context : _output_vexpr_ctxs) { |
270 | 378 | types.emplace_back(context->root()->data_type()); |
271 | 378 | } |
272 | 19 | std::vector<std::string> names = _column_names; |
273 | 19 | if (!_parquet_schemas.empty()) { |
274 | 0 | names.clear(); |
275 | 0 | names.reserve(_parquet_schemas.size()); |
276 | 0 | for (const auto& schema : _parquet_schemas) { |
277 | 0 | names.emplace_back(schema.schema_column_name); |
278 | 0 | } |
279 | 0 | } |
280 | 19 | _arrow_block_convertor = |
281 | 19 | _create_arrow_block_convertor(std::move(types), std::move(names), _timezone, |
282 | 19 | _timezone_obj, _parquet_options.enable_int96_timestamps); |
283 | 19 | RETURN_IF_ERROR(_arrow_block_convertor->init()); |
284 | 19 | try { |
285 | 19 | RETURN_DORIS_STATUS_IF_ERROR(_open_file_writer()); |
286 | 19 | } catch (const ::parquet::ParquetStatusException& e) { |
287 | 0 | LOG(WARNING) << "parquet file writer open error: " << e.what(); |
288 | 0 | return Status::InternalError("parquet file writer open error: {}", e.what()); |
289 | 0 | } |
290 | 19 | if (_writer == nullptr) { |
291 | 0 | return Status::InternalError("Failed to create file writer"); |
292 | 0 | } |
293 | 19 | return Status::OK(); |
294 | 19 | } |
295 | | |
296 | 5 | int64_t VParquetWriter::written_len() { |
297 | 5 | return _outstream->get_written_len(); |
298 | 5 | } |
299 | | |
300 | 19 | Status VParquetWriter::close() { |
301 | 19 | try { |
302 | 19 | if (_writer != nullptr) { |
303 | 19 | RETURN_DORIS_STATUS_IF_ERROR(_writer->Close()); |
304 | 19 | } |
305 | 19 | RETURN_DORIS_STATUS_IF_ERROR(_outstream->Close()); |
306 | | |
307 | 19 | } catch (const std::exception& e) { |
308 | 0 | LOG(WARNING) << "Parquet writer close error: " << e.what(); |
309 | 0 | return Status::IOError(e.what()); |
310 | 0 | } |
311 | | |
312 | 19 | return Status::OK(); |
313 | 19 | } |
314 | | |
315 | | } // namespace doris |