Coverage Report

Created: 2026-09-29 23:09

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
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