Coverage Report

Created: 2026-09-29 23:09

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/exec/sink/writer/vhive_partition_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 <gen_cpp/DataSinks_types.h>
21
22
#include <optional>
23
24
#include "core/column/column.h"
25
#include "exprs/vexpr_fwd.h"
26
#include "format/transformer/vfile_format_transformer.h"
27
#include "io/fs/file_writer.h"
28
29
namespace doris {
30
namespace io {
31
class FileSystem;
32
}
33
34
class ObjectPool;
35
class RuntimeState;
36
class RuntimeProfile;
37
class THiveColumn;
38
39
class Block;
40
class VFileFormatTransformer;
41
42
class VHivePartitionWriter {
43
public:
44
    struct WriteInfo {
45
        std::string write_path;
46
        std::string original_write_path;
47
        std::string target_path;
48
        TFileType::type file_type;
49
        std::vector<TNetworkAddress> broker_addresses;
50
    };
51
52
    VHivePartitionWriter(const TDataSink& t_sink, std::string partition_name,
53
                         TUpdateMode::type update_mode,
54
                         const VExprContextSPtrs& write_output_expr_ctxs,
55
                         std::vector<std::string> write_column_names, WriteInfo write_info,
56
                         std::string file_name, int file_name_index,
57
                         TFileFormatType::type file_format_type,
58
                         TFileCompressType::type hive_compress_type,
59
                         const THiveSerDeProperties* hive_serde_properties,
60
                         const std::map<std::string, std::string>& hadoop_conf);
61
62
0
    Status init_properties(ObjectPool* pool) { return Status::OK(); }
63
64
    Status open(RuntimeState* state, RuntimeProfile* profile);
65
66
    Status write(Block& block);
67
68
    Status close(const Status& status);
69
70
0
    inline const std::string& file_name() const { return _file_name; }
71
72
0
    inline int file_name_index() const { return _file_name_index; }
73
74
0
    inline size_t written_len() { return _file_format_transformer->written_len(); }
75
76
private:
77
    std::string _get_target_file_name();
78
79
private:
80
    THivePartitionUpdate _build_partition_update();
81
    bool _build_s3_mpu_pending_upload(TS3MPUPendingUpload* pending_upload);
82
    void _add_s3_mpu_pending_upload_for_rollback();
83
84
    std::string _get_file_extension(TFileFormatType::type file_format_type,
85
                                    TFileCompressType::type write_compress_type);
86
87
    std::string _path;
88
89
    std::string _partition_name;
90
91
    TUpdateMode::type _update_mode;
92
93
    size_t _row_count = 0;
94
95
    const VExprContextSPtrs& _write_output_expr_ctxs;
96
97
    std::vector<std::string> _write_column_names;
98
99
    WriteInfo _write_info;
100
    std::string _file_name;
101
    int _file_name_index;
102
    TFileFormatType::type _file_format_type;
103
    TFileCompressType::type _hive_compress_type;
104
    const THiveSerDeProperties* _hive_serde_properties;
105
    const std::map<std::string, std::string>& _hadoop_conf;
106
    bool _supports_deferred_azure_multipart = false;
107
    std::optional<std::string> _hive_parquet_time_zone;
108
109
    std::shared_ptr<io::FileSystem> _fs = nullptr;
110
111
    // If the result file format is plain text, like CSV, this _file_writer is owned by this FileResultWriter.
112
    // If the result file format is Parquet, this _file_writer is owned by _parquet_writer.
113
    std::unique_ptr<doris::io::FileWriter> _file_writer = nullptr;
114
    // convert block to parquet/orc/csv format
115
    std::unique_ptr<VFileFormatTransformer> _file_format_transformer = nullptr;
116
117
    RuntimeState* _state;
118
};
119
} // namespace doris