be/src/format_v2/table/hudi_reader.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_v2/table/hudi_reader.h" |
19 | | |
20 | | #include <utility> |
21 | | |
22 | | #include "exprs/vexpr_context.h" |
23 | | #include "format_v2/column_mapper.h" |
24 | | #include "format_v2/jni/hudi_jni_reader.h" |
25 | | #include "format_v2/table/schema_history_util.h" |
26 | | #include "gen_cpp/PlanNodes_types.h" |
27 | | |
28 | | namespace doris::format::hudi { |
29 | | |
30 | 7 | Status HudiReader::prepare_split(const format::SplitReadOptions& options) { |
31 | 7 | { |
32 | | // Derived schema selection is additive to, not nested around, the common base timers. |
33 | 7 | SCOPED_TIMER(_profile.total_timer); |
34 | 7 | SCOPED_TIMER(_profile.prepare_split_timer); |
35 | 7 | _split_schema_id = -1; |
36 | 7 | if (options.current_range.__isset.table_format_params && |
37 | 7 | options.current_range.table_format_params.__isset.hudi_params && |
38 | 7 | options.current_range.table_format_params.hudi_params.__isset.schema_id) { |
39 | 5 | _split_schema_id = options.current_range.table_format_params.hudi_params.schema_id; |
40 | 5 | } |
41 | 7 | } |
42 | 7 | RETURN_IF_ERROR(format::TableReader::prepare_split(options)); |
43 | 7 | SCOPED_TIMER(_profile.total_timer); |
44 | 7 | SCOPED_TIMER(_profile.prepare_split_timer); |
45 | 7 | if (current_split_pruned()) { |
46 | 0 | return Status::OK(); |
47 | 0 | } |
48 | | // This native reader only receives Hudi base-file splits that do not require merging delta |
49 | | // logs. Hudi creates a new versioned base file for updates and compaction instead of modifying |
50 | | // an existing Parquet/ORC base file in place, so missing mtime does not make its page-cache key |
51 | | // ambiguous. Merge-on-read log files remain on the JNI path because they may be appended and |
52 | | // must never be marked immutable here. |
53 | 7 | mark_current_data_file_immutable(); |
54 | 7 | return Status::OK(); |
55 | 7 | } |
56 | | |
57 | 8 | format::TableColumnMappingMode HudiReader::mapping_mode() const { |
58 | 8 | return format::can_map_by_history_schema(_scan_params, _split_schema_id) |
59 | 8 | ? format::TableColumnMappingMode::BY_FIELD_ID |
60 | 8 | : format::TableColumnMappingMode::BY_NAME; |
61 | 8 | } |
62 | | |
63 | 3 | Status HudiReader::annotate_file_schema(std::vector<format::ColumnDefinition>* file_schema) { |
64 | 3 | DORIS_CHECK(file_schema != nullptr); |
65 | 3 | if (mapping_mode() != format::TableColumnMappingMode::BY_FIELD_ID) { |
66 | 2 | return Status::OK(); |
67 | 2 | } |
68 | 1 | return format::annotate_file_schema_from_history(_scan_params, _split_schema_id, file_schema); |
69 | 3 | } |
70 | | |
71 | 2 | Status HudiHybridReader::init(format::TableReadOptions&& options) { |
72 | 2 | return format::TableReader::init(std::move(options)); |
73 | 2 | } |
74 | | |
75 | 3 | Status HudiHybridReader::prepare_split(const format::SplitReadOptions& options) { |
76 | | // A newly selected child initializes against the same scanner profile. Keep hybrid dispatch |
77 | | // outside those shared counters so first-split initialization is counted exactly once. |
78 | 3 | RETURN_IF_ERROR(_ensure_current_split_reader(options)); |
79 | 3 | DORIS_CHECK(_current_split_reader != nullptr); |
80 | 3 | return _current_split_reader->prepare_split(options); |
81 | 3 | } |
82 | | |
83 | 1 | Status HudiHybridReader::get_block(Block* block, bool* eos) { |
84 | 1 | DORIS_CHECK(_current_split_reader != nullptr); |
85 | 1 | return _current_split_reader->get_block(block, eos); |
86 | 1 | } |
87 | | |
88 | 0 | bool HudiHybridReader::current_split_pruned() const { |
89 | 0 | DORIS_CHECK(_current_split_reader != nullptr); |
90 | 0 | return _current_split_reader->current_split_pruned(); |
91 | 0 | } |
92 | | |
93 | 1 | bool HudiHybridReader::current_split_uses_metadata_count() const { |
94 | 1 | DORIS_CHECK(_current_split_reader != nullptr); |
95 | 1 | return _current_split_reader->current_split_uses_metadata_count(); |
96 | 1 | } |
97 | | |
98 | 0 | Status HudiHybridReader::abort_split() { |
99 | 0 | DORIS_CHECK(_current_split_reader != nullptr); |
100 | 0 | return _current_split_reader->abort_split(); |
101 | 0 | } |
102 | | |
103 | 1 | Status HudiHybridReader::close() { |
104 | 1 | Status close_status = Status::OK(); |
105 | 1 | if (_native_reader != nullptr) { |
106 | 1 | close_status = _native_reader->close(); |
107 | 1 | } |
108 | 1 | if (_jni_reader != nullptr) { |
109 | 0 | auto status = _jni_reader->close(); |
110 | 0 | if (!status.ok() && close_status.ok()) { |
111 | 0 | close_status = std::move(status); |
112 | 0 | } |
113 | 0 | } |
114 | 1 | _current_split_reader = nullptr; |
115 | 1 | return close_status; |
116 | 1 | } |
117 | | |
118 | 1 | void HudiHybridReader::set_batch_size(size_t batch_size) { |
119 | 1 | format::TableReader::set_batch_size(batch_size); |
120 | 1 | if (_native_reader != nullptr) { |
121 | 1 | _native_reader->set_batch_size(_batch_size); |
122 | 1 | } |
123 | 1 | if (_jni_reader != nullptr) { |
124 | 1 | _jni_reader->set_batch_size(_batch_size); |
125 | 1 | } |
126 | 1 | } |
127 | | |
128 | 3 | Status HudiHybridReader::_ensure_current_split_reader(const format::SplitReadOptions& options) { |
129 | 3 | DORIS_CHECK(_scan_params != nullptr); |
130 | 3 | if (_is_jni_split(*_scan_params, options.current_range)) { |
131 | 1 | if (_jni_reader == nullptr) { |
132 | 1 | #ifdef BE_TEST |
133 | 1 | if (_test_jni_reader_factory) { |
134 | 1 | _jni_reader = _test_jni_reader_factory(); |
135 | 1 | } else { |
136 | 0 | _jni_reader = std::make_unique<format::hudi::HudiJniReader>(); |
137 | 0 | } |
138 | | #else |
139 | | _jni_reader = std::make_unique<format::hudi::HudiJniReader>(); |
140 | | #endif |
141 | 1 | RETURN_IF_ERROR(_init_child_reader(_jni_reader.get(), format::FileFormat::JNI)); |
142 | 1 | } |
143 | 1 | _current_split_reader = _jni_reader.get(); |
144 | 2 | } else { |
145 | 2 | format::FileFormat file_format; |
146 | 2 | RETURN_IF_ERROR(_to_file_format(*_scan_params, options.current_range, &file_format)); |
147 | 2 | if (_native_reader == nullptr) { |
148 | 2 | #ifdef BE_TEST |
149 | 2 | if (_test_native_reader_factory) { |
150 | 1 | _native_reader = _test_native_reader_factory(); |
151 | 1 | } else { |
152 | 1 | _native_reader = format::hudi::HudiReader::create_unique(); |
153 | 1 | } |
154 | | #else |
155 | | _native_reader = format::hudi::HudiReader::create_unique(); |
156 | | #endif |
157 | 2 | RETURN_IF_ERROR(_init_child_reader(_native_reader.get(), file_format)); |
158 | 2 | } |
159 | 2 | _current_split_reader = _native_reader.get(); |
160 | 2 | } |
161 | 3 | return Status::OK(); |
162 | 3 | } |
163 | | |
164 | | Status HudiHybridReader::_init_child_reader(format::TableReader* reader, |
165 | 3 | format::FileFormat file_format) { |
166 | 3 | DORIS_CHECK(reader != nullptr); |
167 | 3 | VExprContextSPtrs conjuncts; |
168 | 3 | RETURN_IF_ERROR(_clone_conjuncts(&conjuncts)); |
169 | 3 | RETURN_IF_ERROR(reader->init({ |
170 | 3 | .projected_columns = _projected_columns, |
171 | 3 | .conjuncts = std::move(conjuncts), |
172 | 3 | .format = file_format, |
173 | 3 | .scan_params = _scan_params, |
174 | 3 | .io_ctx = _io_ctx, |
175 | 3 | .runtime_state = _runtime_state, |
176 | 3 | .scanner_profile = _scanner_profile, |
177 | 3 | .push_down_agg_type = _push_down_agg_type, |
178 | 3 | .push_down_count_columns = _push_down_count_columns, |
179 | 3 | .condition_cache_digest = _condition_cache_digest, |
180 | 3 | })); |
181 | | // Zero means no adaptive prediction has been produced yet. Preserve the child's normal |
182 | | // runtime default until FileScannerV2 supplies the first positive prediction. |
183 | 3 | if (_batch_size > 0) { |
184 | 0 | reader->set_batch_size(_batch_size); |
185 | 0 | } |
186 | 3 | return Status::OK(); |
187 | 3 | } |
188 | | |
189 | 3 | Status HudiHybridReader::_clone_conjuncts(VExprContextSPtrs* conjuncts) const { |
190 | 3 | DORIS_CHECK(conjuncts != nullptr); |
191 | 3 | conjuncts->clear(); |
192 | 3 | conjuncts->reserve(_conjuncts.size()); |
193 | 3 | for (const auto& conjunct : _conjuncts) { |
194 | 0 | VExprSPtr root; |
195 | 0 | RETURN_IF_ERROR(format::clone_table_expr_tree(conjunct->root(), &root)); |
196 | 0 | conjuncts->push_back(VExprContext::create_shared(std::move(root))); |
197 | 0 | } |
198 | 3 | return Status::OK(); |
199 | 3 | } |
200 | | |
201 | | TFileFormatType::type HudiHybridReader::_range_format_type(const TFileScanRangeParams& params, |
202 | 5 | const TFileRangeDesc& range) { |
203 | 5 | return range.__isset.format_type ? range.format_type : params.format_type; |
204 | 5 | } |
205 | | |
206 | | bool HudiHybridReader::_is_jni_split(const TFileScanRangeParams& params, |
207 | 3 | const TFileRangeDesc& range) { |
208 | 3 | return _range_format_type(params, range) == TFileFormatType::FORMAT_JNI; |
209 | 3 | } |
210 | | |
211 | | Status HudiHybridReader::_to_file_format(const TFileScanRangeParams& params, |
212 | | const TFileRangeDesc& range, |
213 | 2 | format::FileFormat* file_format) { |
214 | 2 | DORIS_CHECK(file_format != nullptr); |
215 | 2 | const auto format_type = _range_format_type(params, range); |
216 | 2 | switch (format_type) { |
217 | 2 | case TFileFormatType::FORMAT_PARQUET: |
218 | 2 | *file_format = format::FileFormat::PARQUET; |
219 | 2 | return Status::OK(); |
220 | 0 | case TFileFormatType::FORMAT_ORC: |
221 | 0 | *file_format = format::FileFormat::ORC; |
222 | 0 | return Status::OK(); |
223 | 0 | default: |
224 | 0 | return Status::NotSupported("Unsupported native Hudi file format {}", |
225 | 0 | to_string(format_type)); |
226 | 2 | } |
227 | 2 | } |
228 | | |
229 | | } // namespace doris::format::hudi |