Coverage Report

Created: 2026-09-29 18:23

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