Coverage Report

Created: 2026-09-25 14:54

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/format/transformer/vparquet_writer.h
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
#pragma once
19
20
#include <arrow/io/interfaces.h>
21
#include <arrow/result.h>
22
#include <arrow/status.h>
23
#include <cctz/time_zone.h>
24
#include <gen_cpp/DataSinks_types.h>
25
#include <parquet/arrow/writer.h>
26
#include <parquet/file_writer.h>
27
#include <parquet/properties.h>
28
#include <parquet/types.h>
29
30
#include <cstdint>
31
32
#include "format/arrow/arrow_block_convertor.h"
33
#include "format/transformer/vfile_format_transformer.h"
34
35
namespace doris {
36
namespace io {
37
class FileWriter;
38
} // namespace io
39
} // namespace doris
40
namespace parquet {
41
namespace schema {
42
class GroupNode;
43
} // namespace schema
44
} // namespace parquet
45
46
namespace doris {
47
48
class ParquetOutputStream : public arrow::io::OutputStream {
49
public:
50
    ParquetOutputStream(doris::io::FileWriter* file_writer);
51
    ParquetOutputStream(doris::io::FileWriter* file_writer, const int64_t& written_len);
52
    ~ParquetOutputStream() override;
53
54
    arrow::Status Write(const void* data, int64_t nbytes) override;
55
    // return the current write position of the stream
56
    arrow::Result<int64_t> Tell() const override;
57
    arrow::Status Close() override;
58
59
0
    bool closed() const override { return _is_closed; }
60
61
    int64_t get_written_len() const;
62
63
    void set_written_len(int64_t written_len);
64
65
private:
66
    doris::io::FileWriter* _file_writer = nullptr; // not owned
67
    int64_t _cur_pos = 0;                          // current write position
68
    bool _is_closed = false;
69
    int64_t _written_len = 0;
70
};
71
72
class ParquetBuildHelper {
73
public:
74
    static void build_compression_type(::parquet::WriterProperties::Builder& builder,
75
                                       const TParquetCompressionType::type& compression_type);
76
77
    static void build_version(::parquet::WriterProperties::Builder& builder,
78
                              const TParquetVersion::type& parquet_version);
79
};
80
81
struct ParquetFileOptions {
82
    TParquetCompressionType::type compression_type;
83
    TParquetVersion::type parquet_version;
84
    bool parquet_disable_dictionary = false;
85
    bool enable_int96_timestamps = false;
86
};
87
88
// Writes Doris blocks as Parquet files, including schema and Arrow conversion.
89
class VParquetWriter : public VFileFormatTransformer {
90
public:
91
    VParquetWriter(RuntimeState* state, doris::io::FileWriter* file_writer,
92
                   const VExprContextSPtrs& output_vexpr_ctxs,
93
                   std::vector<std::string> column_names, bool output_object_data,
94
                   const ParquetFileOptions& parquet_options);
95
96
    VParquetWriter(RuntimeState* state, doris::io::FileWriter* file_writer,
97
                   const VExprContextSPtrs& output_vexpr_ctxs,
98
                   std::vector<TParquetSchema> parquet_schemas, bool output_object_data,
99
                   const ParquetFileOptions& parquet_options);
100
101
3
    ~VParquetWriter() override = default;
102
103
    Status open() override;
104
105
    Status write(const Block& block) override;
106
107
    Status close() override;
108
109
    int64_t written_len() override;
110
111
protected:
112
    // Construct the schema and column bindings together for each writer instance.
113
    virtual std::unique_ptr<ArrowBlockConvertor> _create_arrow_block_convertor(
114
            DataTypes types, std::vector<std::string> names, const std::string& timezone_name,
115
            const cctz::time_zone& timezone) const;
116
1
    std::shared_ptr<::parquet::FileMetaData> _file_metadata() const { return _writer->metadata(); }
117
118
private:
119
    std::unique_ptr<ArrowBlockConvertor> _arrow_block_convertor;
120
    Status _parse_properties();
121
    arrow::Status _open_file_writer();
122
123
    std::shared_ptr<ParquetOutputStream> _outstream;
124
    std::shared_ptr<::parquet::WriterProperties> _parquet_writer_properties;
125
    std::shared_ptr<::parquet::ArrowWriterProperties> _arrow_properties;
126
    std::unique_ptr<::parquet::arrow::FileWriter> _writer;
127
128
    std::vector<std::string> _column_names;
129
    std::vector<TParquetSchema> _parquet_schemas;
130
    const ParquetFileOptions _parquet_options;
131
    std::string _timezone;
132
    cctz::time_zone _timezone_obj;
133
    uint64_t _write_size = 0;
134
};
135
136
} // namespace doris