be/src/format_v2/table_reader.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 <bvar/status.h> |
21 | | |
22 | | #include <algorithm> |
23 | | #include <exception> |
24 | | #include <map> |
25 | | #include <memory> |
26 | | #include <optional> |
27 | | #include <string> |
28 | | #include <string_view> |
29 | | #include <utility> |
30 | | #include <vector> |
31 | | |
32 | | #include "common/cast_set.h" |
33 | | #include "common/exception.h" |
34 | | #include "common/logging.h" |
35 | | #include "common/status.h" |
36 | | #include "core/assert_cast.h" |
37 | | #include "core/block/block.h" |
38 | | #include "core/column/column_array.h" |
39 | | #include "core/column/column_const.h" |
40 | | #include "core/column/column_map.h" |
41 | | #include "core/column/column_nullable.h" |
42 | | #include "core/column/column_struct.h" |
43 | | #include "core/column/column_vector.h" |
44 | | #include "core/data_type/data_type.h" |
45 | | #include "core/data_type/data_type_array.h" |
46 | | #include "core/data_type/data_type_map.h" |
47 | | #include "core/data_type/data_type_nullable.h" |
48 | | #include "core/data_type/data_type_number.h" |
49 | | #include "core/data_type/data_type_string.h" |
50 | | #include "core/data_type/data_type_struct.h" |
51 | | #include "core/field.h" |
52 | | #include "exec/common/stringop_substring.h" |
53 | | #include "exprs/vexpr.h" |
54 | | #include "exprs/vexpr_context.h" |
55 | | #include "exprs/vexpr_fwd.h" |
56 | | #include "exprs/vslot_ref.h" |
57 | | #include "format/table/deletion_vector.h" |
58 | | #include "format_v2/column_data.h" |
59 | | #include "format_v2/column_mapper.h" |
60 | | #include "format_v2/expr/cast.h" |
61 | | #include "format_v2/expr/delete_predicate.h" |
62 | | #include "format_v2/file_reader.h" |
63 | | #include "format_v2/parquet/reader/column_reader.h" |
64 | | #include "format_v2/schema_projection.h" |
65 | | #include "gen_cpp/PlanNodes_types.h" |
66 | | #include "io/io_common.h" |
67 | | #include "runtime/descriptors.h" |
68 | | #include "storage/segment/condition_cache.h" |
69 | | |
70 | | namespace doris { |
71 | | class Block; |
72 | | struct DeleteFileDesc; |
73 | | class RuntimeState; |
74 | | } // namespace doris |
75 | | |
76 | | namespace doris::format { |
77 | | |
78 | | using DeleteRows = std::vector<int64_t>; |
79 | | |
80 | | // Row-level predicates on table/global schema. They are rewritten to file-local expressions when |
81 | | // possible, and remain the source of row-level filtering after localization. |
82 | | struct TableFilter { |
83 | | VExprContextSPtr conjunct; |
84 | | std::vector<GlobalIndex> global_indices; |
85 | | }; |
86 | | |
87 | | struct ScanTask { |
88 | 128k | virtual ~ScanTask() = default; |
89 | | |
90 | | std::unique_ptr<io::FileDescription> data_file; |
91 | | }; |
92 | | |
93 | | struct ProjectedColumnBuildContext { |
94 | | const TFileScanRangeParams* scan_params = nullptr; |
95 | | const TFileRangeDesc* range = nullptr; |
96 | | RuntimeState* runtime_state = nullptr; |
97 | | const SlotDescriptor* slot_desc = nullptr; |
98 | | std::optional<ColumnDefinition> schema_column = std::nullopt; |
99 | | size_t next_file_column_idx = 0; |
100 | | }; |
101 | | |
102 | | struct ReadProfile { |
103 | | RuntimeProfile::Counter* total_timer = nullptr; |
104 | | RuntimeProfile::Counter* init_timer = nullptr; |
105 | | RuntimeProfile::Counter* num_delete_files = nullptr; |
106 | | RuntimeProfile::Counter* num_delete_rows = nullptr; |
107 | | RuntimeProfile::Counter* parse_delete_file_time = nullptr; |
108 | | RuntimeProfile::Counter* decoded_dv_cache_hit_count = nullptr; |
109 | | RuntimeProfile::Counter* decoded_dv_cache_miss_count = nullptr; |
110 | | RuntimeProfile::Counter* dv_file_cache_hit_count = nullptr; |
111 | | RuntimeProfile::Counter* dv_file_cache_miss_count = nullptr; |
112 | | RuntimeProfile::Counter* dv_file_cache_peer_read_count = nullptr; |
113 | | RuntimeProfile::Counter* exec_timer = nullptr; |
114 | | RuntimeProfile::Counter* prepare_split_timer = nullptr; |
115 | | RuntimeProfile::Counter* finalize_timer = nullptr; |
116 | | RuntimeProfile::Counter* create_reader_timer = nullptr; |
117 | | RuntimeProfile::Counter* pushdown_agg_timer = nullptr; |
118 | | RuntimeProfile::Counter* open_reader_timer = nullptr; |
119 | | RuntimeProfile::Counter* refresh_conjuncts_timer = nullptr; |
120 | | RuntimeProfile::Counter* runtime_filter_partition_prune_timer = nullptr; |
121 | | RuntimeProfile::Counter* runtime_filter_partition_pruned_range_counter = nullptr; |
122 | | RuntimeProfile::Counter* close_timer = nullptr; |
123 | | RuntimeProfile::Counter* file_reader_total_timer = nullptr; |
124 | | RuntimeProfile::Counter* file_reader_init_timer = nullptr; |
125 | | RuntimeProfile::Counter* file_reader_schema_timer = nullptr; |
126 | | RuntimeProfile::Counter* file_reader_mapper_timer = nullptr; |
127 | | RuntimeProfile::Counter* file_reader_open_timer = nullptr; |
128 | | RuntimeProfile::Counter* file_reader_refresh_timer = nullptr; |
129 | | RuntimeProfile::Counter* file_reader_get_block_timer = nullptr; |
130 | | RuntimeProfile::Counter* file_reader_aggregate_timer = nullptr; |
131 | | RuntimeProfile::Counter* file_reader_close_timer = nullptr; |
132 | | }; |
133 | | |
134 | | struct TableReadOptions { |
135 | | // Columns need to be read from file and output by table reader. They are all in table/global |
136 | | // schema semantics. |
137 | | const std::vector<ColumnDefinition> projected_columns; |
138 | | // All complex conjuncts from scan operator |
139 | | const VExprContextSPtrs conjuncts; |
140 | | // File format of the underlying data files, needed for reader initialization and reader-level |
141 | | // filter pushdown. |
142 | | const FileFormat format; |
143 | | TFileScanRangeParams* scan_params; |
144 | | std::shared_ptr<io::IOContext> io_ctx; |
145 | | RuntimeState* runtime_state; |
146 | | RuntimeProfile* scanner_profile; |
147 | | // File formats without complete self-describing metadata, such as CSV, Text, and JSON, need |
148 | | // the FE-planned physical file slots to build their file-local schema and deserialize values. |
149 | | const std::vector<SlotDescriptor*>* file_slot_descs = nullptr; |
150 | | // Push-down aggregate type. |
151 | | const TPushAggOp::type push_down_agg_type = TPushAggOp::type::NONE; |
152 | | // Table/global indices of explicit COUNT arguments. nullopt means an old FE did not send the |
153 | | // semantic argument field, while an explicit empty vector means COUNT(*)/COUNT(1). Keeping |
154 | | // those states separate prevents a rolling-upgrade plan from being reinterpreted by a new BE. |
155 | | const std::optional<std::vector<GlobalIndex>> push_down_count_columns = std::nullopt; |
156 | | // Initial digest of predicates available during scanner open. Scanner-driven splits override it |
157 | | // with SplitReadOptions::condition_cache_digest after collecting late-arrival runtime filters. |
158 | | // A zero digest disables condition cache. |
159 | | uint64_t condition_cache_digest = 0; |
160 | | }; |
161 | | |
162 | | struct SplitReadOptions { |
163 | | // Split-level information for reader initialization, which may include file path, partition values, delete file info, etc. The content is table format specific and opaque to table reader base class; it's the responsibility of the concrete table reader implementation to parse necessary information for reader initialization and filter pushdown. |
164 | | std::map<std::string, Field> partition_values; |
165 | | // Latest scanner conjuncts rewritten to table/global column indices. Runtime filters may |
166 | | // arrive after TableReader::init(), so scanner-driven splits replace the initial snapshot. |
167 | | // nullopt preserves the initial snapshot for standalone TableReader callers. |
168 | | std::optional<VExprContextSPtrs> conjuncts = std::nullopt; |
169 | | // Independent clones used for partition pruning because evaluation prepares and opens them |
170 | | // against a synthetic partition block before the file reader opens its row-level conjuncts. |
171 | | VExprContextSPtrs partition_prune_conjuncts; |
172 | | // Table-level COUNT may emit one metadata-derived batch and resume on a later scheduler turn. |
173 | | // It is safe only after every runtime filter assigned to the scanner has arrived; otherwise a |
174 | | // filter could arrive after synthetic rows have already been returned and those rows cannot be |
175 | | // retracted. Standalone TableReader callers have no scanner runtime-filter lifecycle. |
176 | | bool all_runtime_filters_applied = true; |
177 | | // Digest for the exact scanner conjunct snapshot attached to this split. FileScannerV2 rebuilds |
178 | | // it after collecting late-arrival RFs, so different RF payloads cannot share a cache entry. A |
179 | | // zero value explicitly disables condition cache for this split. |
180 | | std::optional<uint64_t> condition_cache_digest; |
181 | | ShardedKVCache* cache = nullptr; |
182 | | TFileRangeDesc current_range; |
183 | | FileFormat current_split_format = FileFormat::PARQUET; |
184 | | std::optional<GlobalRowIdContext> global_rowid_context; |
185 | | }; |
186 | | |
187 | | // Base class for table-level readers. |
188 | | // This layer owns common table-level orchestration, such as split iteration, dynamic partition |
189 | | // pruning, delete handling and conversion from file-local blocks to table-schema blocks. Concrete |
190 | | // table-format readers only need to provide format-specific hooks for opening readers and parsing |
191 | | // split metadata. |
192 | | class TableReader { |
193 | | public: |
194 | 60.3k | virtual ~TableReader() = default; |
195 | | |
196 | | // Initialize common runtime options for the table reader. Subclasses may call this from their |
197 | | // own init(options); table-format schema and split metadata are provided later per split. |
198 | | virtual Status init(TableReadOptions&& options); |
199 | | |
200 | | // FileScannerV2 adjusts this before each get_block() using an adaptive bytes-per-row estimate. |
201 | | // Store it here as well as forwarding to the current reader so newly opened split readers start |
202 | | // with the latest predicted batch size. |
203 | 493k | virtual void set_batch_size(size_t batch_size) { |
204 | 493k | _batch_size = std::max<size_t>(1, batch_size); |
205 | 493k | if (_data_reader.reader != nullptr) { |
206 | 108k | _data_reader.reader->set_batch_size(_batch_size); |
207 | 108k | } |
208 | 493k | } |
209 | | |
210 | | #ifdef BE_TEST |
211 | | size_t TEST_batch_size() const { return _batch_size; } |
212 | | void TEST_set_condition_cache_hit_count(int64_t hits) { _condition_cache_hit_count = hits; } |
213 | | bool TEST_current_data_file_is_immutable() const { |
214 | | DORIS_CHECK(_current_task != nullptr); |
215 | | DORIS_CHECK(_current_task->data_file != nullptr); |
216 | | DORIS_CHECK(_current_file_description.has_value()); |
217 | | DORIS_CHECK(_current_task->data_file->is_immutable == |
218 | | _current_file_description->is_immutable); |
219 | | return _current_task->data_file->is_immutable; |
220 | | } |
221 | | #endif |
222 | | |
223 | | // Prepare for reading a new split/task. |
224 | | // 1. Pass a new split/task to reader, which will be used in subsequent open_reader() to initialize the underlying file reader. |
225 | | // 2. Parse delete predicates from split/task information, which will be used for later dynamic filtering and delete handling. |
226 | | virtual Status prepare_split(const SplitReadOptions& options); |
227 | | |
228 | | // Refresh row-level predicates for an already prepared split. Physical readers that support |
229 | | // this operation decide the safe boundary at which the new immutable request becomes active. |
230 | | virtual Status refresh_conjuncts(VExprContextSPtrs conjuncts); |
231 | | |
232 | 224k | virtual bool current_split_pruned() const { return _current_split_pruned; } |
233 | 364k | virtual bool current_split_uses_metadata_count() const { |
234 | 364k | return _current_split_uses_metadata_count; |
235 | 364k | } |
236 | | |
237 | | // Discard the active split after the caller decides an error is ignorable, for example a |
238 | | // stale external-table file listing that returns NOT_FOUND. The next prepare_split() must start |
239 | | // with no concrete reader or split-local state left from the failed split. |
240 | 1 | virtual Status abort_split() { |
241 | | // Ignored open failures still spend time closing partially initialized readers. Include |
242 | | // that recovery path in the common lifecycle profile so NOT_FOUND cannot become invisible. |
243 | 1 | SCOPED_TIMER(_profile.total_timer); |
244 | 1 | SCOPED_TIMER(_profile.close_timer); |
245 | 1 | if (_data_reader.reader != nullptr) { |
246 | 1 | RETURN_IF_ERROR(close_current_reader()); |
247 | 1 | } else { |
248 | 0 | _current_task.reset(); |
249 | 0 | _current_file_description.reset(); |
250 | 0 | } |
251 | 1 | _delete_rows = nullptr; |
252 | 1 | _remaining_table_level_count = -1; |
253 | 1 | _remaining_file_level_count = -1; |
254 | 1 | _current_split_uses_metadata_count = false; |
255 | 1 | _current_split_pruned = false; |
256 | 1 | return Status::OK(); |
257 | 1 | } |
258 | | |
259 | | // Public entry point for reading a table-schema block. The base class opens the current reader, |
260 | | // advances across EOF, and closes exhausted readers. Subclasses provide protected hooks for |
261 | | // table-format-specific behavior. |
262 | 232k | virtual Status get_block(Block* block, bool* eos) { |
263 | 232k | SCOPED_TIMER(_profile.total_timer); |
264 | 232k | SCOPED_TIMER(_profile.exec_timer); |
265 | 232k | DORIS_CHECK(block->columns() == _projected_columns.size()); |
266 | 232k | block->clear_column_data(_projected_columns.size()); |
267 | | |
268 | 363k | while (true) { |
269 | 363k | if (*eos) { |
270 | 0 | return Status::OK(); |
271 | 0 | } |
272 | 363k | if (_io_ctx != nullptr && _io_ctx->should_stop) { |
273 | 10 | *eos = true; |
274 | 10 | return Status::OK(); |
275 | 10 | } |
276 | 363k | if (!_data_reader.reader) { |
277 | 243k | if (_is_table_level_count_active()) { |
278 | 314 | RETURN_IF_ERROR(_read_table_level_count(block, eos)); |
279 | 314 | return Status::OK(); |
280 | 314 | } |
281 | 243k | if (_is_file_level_count_active()) { |
282 | 2.93k | RETURN_IF_ERROR(_read_file_level_count(block, eos)); |
283 | 2.93k | return Status::OK(); |
284 | 2.93k | } |
285 | 240k | RETURN_IF_ERROR(create_next_reader(eos)); |
286 | 240k | if (!_data_reader.reader) { |
287 | 119k | DCHECK(*eos); |
288 | 119k | return Status::OK(); |
289 | 119k | } |
290 | 240k | } |
291 | | |
292 | | // Materialize a reduced row set for upper aggregate operators when aggregate |
293 | | // pushdown can be applied. This is not the final aggregate result: COUNT emits |
294 | | // `count` default rows for the upper COUNT(*), and MIN/MAX emits two rows containing |
295 | | // file-level min/max values for the upper MIN/MAX. |
296 | 240k | if (!_aggregate_pushdown_tried) { |
297 | 121k | SCOPED_TIMER(_profile.pushdown_agg_timer); |
298 | 121k | bool pushed_down = false; |
299 | 121k | const auto status = _try_materialize_aggregate_pushdown_rows(block, &pushed_down); |
300 | 121k | if (!status.ok()) { |
301 | 1 | if (_io_ctx != nullptr && _io_ctx->should_stop && |
302 | 1 | status.is<ErrorCode::END_OF_FILE>()) { |
303 | 1 | *eos = true; |
304 | 1 | return Status::OK(); |
305 | 1 | } |
306 | 0 | return status; |
307 | 1 | } |
308 | 121k | if (pushed_down) { |
309 | 1.49k | return Status::OK(); |
310 | 1.49k | } |
311 | 121k | } |
312 | | |
313 | 239k | bool current_eof = false; |
314 | 239k | _data_reader.block_template.clear_column_data( |
315 | 239k | cast_set<int64_t>(_data_reader.file_block_layout.size())); |
316 | 239k | size_t current_rows = 0; |
317 | 239k | { |
318 | 239k | SCOPED_TIMER(_profile.file_reader_total_timer); |
319 | 239k | SCOPED_TIMER(_profile.file_reader_get_block_timer); |
320 | 239k | RETURN_IF_ERROR(_data_reader.reader->get_block(&_data_reader.block_template, |
321 | 239k | ¤t_rows, ¤t_eof)); |
322 | 239k | } |
323 | 239k | const bool stopped_during_read = _io_ctx != nullptr && _io_ctx->should_stop; |
324 | 239k | if (current_rows == 0) { |
325 | 130k | if (current_eof) { |
326 | 119k | _current_reader_reached_eof = !stopped_during_read; |
327 | 119k | RETURN_IF_ERROR(close_current_reader()); |
328 | 119k | } |
329 | 130k | continue; |
330 | 130k | } |
331 | 239k | DCHECK_EQ(_data_reader.block_template.columns(), _data_reader.file_block_layout.size()) |
332 | 0 | << _data_reader.block_template.dump_structure(); |
333 | 108k | #ifndef NDEBUG |
334 | 108k | RETURN_IF_ERROR(_check_file_block_columns("after file reader get_block", current_rows)); |
335 | 108k | #endif |
336 | 108k | DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size()); |
337 | 108k | RETURN_IF_ERROR(finalize_chunk(block, current_rows)); |
338 | 108k | #ifndef NDEBUG |
339 | 108k | RETURN_IF_ERROR( |
340 | 108k | _check_table_block_columns("after finalize_chunk", block, current_rows)); |
341 | 108k | #endif |
342 | 108k | if (current_eof) { |
343 | 20 | _current_reader_reached_eof = !stopped_during_read; |
344 | 20 | RETURN_IF_ERROR(close_current_reader()); |
345 | 20 | } |
346 | 108k | return Status::OK(); |
347 | 108k | } |
348 | 232k | } |
349 | | |
350 | | // Close the table reader and the currently active file reader. Subclasses that hold additional |
351 | | // table-format resources should override this and call TableReader::close() first. |
352 | 54.0k | virtual Status close() { |
353 | 54.0k | SCOPED_TIMER(_profile.total_timer); |
354 | 54.0k | SCOPED_TIMER(_profile.close_timer); |
355 | 54.0k | if (_data_reader.reader) { |
356 | 665 | RETURN_IF_ERROR(close_current_reader()); |
357 | 665 | } |
358 | 54.0k | _current_task.reset(); |
359 | 54.0k | _current_file_description.reset(); |
360 | 54.0k | _remaining_table_level_count = -1; |
361 | 54.0k | _remaining_file_level_count = -1; |
362 | 54.0k | _current_split_uses_metadata_count = false; |
363 | 54.0k | return Status::OK(); |
364 | 54.0k | } |
365 | | |
366 | 106k | virtual int64_t condition_cache_hit_count() const { return _condition_cache_hit_count; } |
367 | | |
368 | | virtual std::string debug_string() const; |
369 | | |
370 | | virtual Status annotate_projected_column(const TFileScanSlotInfo& slot_info, |
371 | | ProjectedColumnBuildContext* context, |
372 | | ColumnDefinition* column) const; |
373 | | |
374 | 34.4k | virtual Status validate_projected_columns(const ProjectedColumnBuildContext& context) const { |
375 | 34.4k | (void)context; |
376 | 34.4k | return Status::OK(); |
377 | 34.4k | } |
378 | | |
379 | | protected: |
380 | | // TableReader keeps the active file description both in the scan task and separately for |
381 | | // creating the physical reader. Table-format readers must update both copies when their |
382 | | // snapshot protocol guarantees that a file path is never overwritten with different bytes. |
383 | | // This guarantee lets readers safely build cache keys without mtime; it must not be used for |
384 | | // ordinary Hive/TVF files whose paths may be overwritten in place. |
385 | 89.7k | void mark_current_data_file_immutable() { |
386 | 89.7k | DORIS_CHECK(_current_task != nullptr); |
387 | 89.7k | DORIS_CHECK(_current_task->data_file != nullptr); |
388 | 89.7k | DORIS_CHECK(_current_file_description.has_value()); |
389 | 89.7k | _current_task->data_file->is_immutable = true; |
390 | 89.7k | _current_file_description->is_immutable = true; |
391 | 89.7k | } |
392 | | |
393 | | std::optional<ColumnDefinition> _find_table_column_by_field_id( |
394 | | int32_t field_id, DataTypePtr type, bool include_historical_schemas) const; |
395 | | std::optional<std::vector<ColumnDefinition>> _find_table_column_path_by_field_id( |
396 | | int32_t field_id, DataTypePtr leaf_type, bool include_historical_schemas) const; |
397 | | std::optional<std::vector<ColumnDefinition>> _find_table_column_identity_path_by_field_id( |
398 | | int32_t field_id, bool include_historical_schemas) const; |
399 | | |
400 | | // Parse deletion vector information from table format specific file description. |
401 | | virtual Status _parse_deletion_vector_file(const TTableFormatFileDesc& t_desc, |
402 | 37.3k | DeleteFileDesc* desc, bool* has_delete_file) { |
403 | 37.3k | *has_delete_file = false; |
404 | 37.3k | return Status::OK(); |
405 | 37.3k | } |
406 | | |
407 | | // Advance to the next reader. This closes the current reader first and then opens the next |
408 | | // concrete reader. Subclasses should not duplicate this loop. |
409 | | Status create_next_reader(bool* eos); |
410 | | virtual Status create_file_reader(std::unique_ptr<FileReader>* reader); |
411 | 6.37k | virtual TableColumnMappingMode mapping_mode() const { return TableColumnMappingMode::BY_NAME; } |
412 | 89.5k | virtual void configure_mapper_options(TableColumnMapperOptions*) const {} |
413 | 38.6k | virtual Status annotate_file_schema(std::vector<ColumnDefinition>* file_schema) { |
414 | 38.6k | DORIS_CHECK(file_schema != nullptr); |
415 | 38.6k | return Status::OK(); |
416 | 38.6k | } |
417 | | |
418 | | // Open the concrete reader for the current split/task and build the file-local scan request. |
419 | 122k | virtual Status open_reader() { |
420 | 122k | SCOPED_TIMER(_profile.open_reader_timer); |
421 | | // 1. Get file schema and create column mapping. |
422 | 122k | std::vector<ColumnDefinition> file_schema; |
423 | 122k | { |
424 | 122k | SCOPED_TIMER(_profile.file_reader_total_timer); |
425 | 122k | SCOPED_TIMER(_profile.file_reader_schema_timer); |
426 | 122k | RETURN_IF_ERROR(_data_reader.reader->get_schema(&file_schema)); |
427 | 122k | } |
428 | | // For Paimon/Hudi, FE can provide field ids through `history_schema_info`. Annotate the |
429 | | // file schema before column mapping when the table format maps columns by field id. |
430 | 122k | RETURN_IF_ERROR(annotate_file_schema(&file_schema)); |
431 | 122k | _data_reader.file_schema = file_schema; |
432 | 122k | _mapper_options.mode = mapping_mode(); |
433 | 122k | configure_mapper_options(&_mapper_options); |
434 | | |
435 | 122k | { |
436 | 122k | SCOPED_TIMER(_profile.file_reader_total_timer); |
437 | 122k | SCOPED_TIMER(_profile.file_reader_mapper_timer); |
438 | 122k | _data_reader.column_mapper = _data_reader.reader->create_column_mapper(_mapper_options); |
439 | 122k | } |
440 | 122k | DORIS_CHECK(_data_reader.column_mapper != nullptr); |
441 | 122k | RETURN_IF_ERROR(_data_reader.column_mapper->create_mapping(_projected_columns, |
442 | 122k | _partition_values, file_schema)); |
443 | 122k | DORIS_CHECK(_data_reader.column_mapper->mappings().size() == _projected_columns.size()); |
444 | | |
445 | | // 2. Build table filters based on conjuncts and column predicates. |
446 | 122k | RETURN_IF_ERROR(_build_table_filters_from_conjuncts()); |
447 | | |
448 | | // 3. Create file scan request based on column mapping and table filters, then open file |
449 | | // reader with the request. File scan request carries row-level expression filters and |
450 | | // file-level pruning hints. Only expression filters decide returned rows. |
451 | 122k | auto file_request = std::make_shared<FileScanRequest>(); |
452 | 122k | RETURN_IF_ERROR(_data_reader.column_mapper->create_scan_request( |
453 | 122k | _table_filters, _projected_columns, file_request.get(), _runtime_state)); |
454 | 122k | bool constant_filter_pruned_split = false; |
455 | 122k | RETURN_IF_ERROR(_evaluate_constant_filters(&constant_filter_pruned_split)); |
456 | 122k | if (constant_filter_pruned_split) { |
457 | 627 | RETURN_IF_ERROR(close_current_reader()); |
458 | 627 | return Status::OK(); |
459 | 627 | } |
460 | | // COUNT(*) has no semantic column argument, but Nereids retains a minimum-width scan slot |
461 | | // so the scan node still has an output tuple. Record only the current non-predicate file |
462 | | // columns before table-format hooks add row-position or equality-delete dependencies. This |
463 | | // marker is independent of aggregate eligibility: with position deletes, for example, |
464 | | // metadata COUNT must fall back to reading rows, but an arbitrary unsupported TIME_MILLIS |
465 | | // placeholder still must not be validated or decoded merely to carry the surviving count. |
466 | | // Pending runtime filters may later target this retained slot, so placeholder values are |
467 | | // safe only after every filter for the split has arrived. |
468 | 121k | if (_push_down_agg_type == TPushAggOp::type::COUNT && |
469 | 121k | _push_down_count_columns.has_value() && _push_down_count_columns->empty() && |
470 | 121k | _all_runtime_filters_applied_for_split) { |
471 | 1.96k | file_request->count_star_placeholder_columns.reserve( |
472 | 1.96k | file_request->non_predicate_columns.size()); |
473 | 1.96k | for (const auto& column : file_request->non_predicate_columns) { |
474 | 1.93k | file_request->count_star_placeholder_columns.push_back(column.column_id()); |
475 | 1.93k | } |
476 | 1.96k | } |
477 | 121k | RETURN_IF_ERROR(customize_file_scan_request(file_request.get())); |
478 | 121k | RETURN_IF_ERROR(_open_local_filter_exprs(*file_request)); |
479 | 121k | _data_reader.file_block_layout.clear(); |
480 | 121k | _data_reader.block_template.clear(); |
481 | 121k | _file_scan_request.reset(); |
482 | 121k | _data_reader.file_block_layout.resize(file_request->local_positions.size()); |
483 | | |
484 | | // 4. Build file block layout from file schema and column mapping. The layout describes |
485 | | // the block returned by file reader before table-column materialization. |
486 | 532k | for (const auto& [file_column_id, block_position] : file_request->local_positions) { |
487 | 532k | DORIS_CHECK(block_position.value() < _data_reader.file_block_layout.size()); |
488 | 532k | const auto* field = _find_column_definition(_data_reader.file_schema, file_column_id); |
489 | 532k | DORIS_CHECK(field != nullptr); |
490 | | |
491 | 532k | ColumnDefinition projected_field; |
492 | 532k | { |
493 | 532k | auto it = std::find_if( |
494 | 532k | file_request->non_predicate_columns.begin(), |
495 | 532k | file_request->non_predicate_columns.end(), |
496 | 8.92M | [&](const LocalColumnIndex& p) { return p.column_id() == file_column_id; }); |
497 | 532k | if (it != file_request->non_predicate_columns.end()) { |
498 | 442k | RETURN_IF_ERROR(project_column_definition(*field, *it, &projected_field)); |
499 | 442k | } |
500 | 532k | } |
501 | 532k | { |
502 | 532k | auto it = std::find_if( |
503 | 532k | file_request->predicate_columns.begin(), |
504 | 532k | file_request->predicate_columns.end(), |
505 | 532k | [&](const LocalColumnIndex& p) { return p.column_id() == file_column_id; }); |
506 | 532k | if (it != file_request->predicate_columns.end()) { |
507 | 90.4k | RETURN_IF_ERROR(project_column_definition(*field, *it, &projected_field)); |
508 | 90.4k | } |
509 | 532k | } |
510 | 532k | _data_reader.file_block_layout[block_position.value()] = { |
511 | 532k | .file_column_id = file_column_id, |
512 | 532k | .name = projected_field.name, |
513 | 532k | .type = projected_field.type, |
514 | 532k | }; |
515 | 532k | DORIS_CHECK(_data_reader.file_block_layout[block_position.value()].type != nullptr); |
516 | 532k | } |
517 | | |
518 | | // 5. Prepare block template from file block layout. The block template stores the block |
519 | | // returned by file reader before table-column materialization. |
520 | 121k | _data_reader.block_template.reserve(_data_reader.file_block_layout.size()); |
521 | 531k | for (const auto& column : _data_reader.file_block_layout) { |
522 | 531k | _data_reader.block_template.insert( |
523 | 531k | {column.type->create_column(), column.type, column.name}); |
524 | 531k | } |
525 | 121k | if (VLOG_DEBUG_IS_ON) { |
526 | 0 | VLOG_DEBUG << "TableReader debug: " << debug_string(); |
527 | 0 | } |
528 | 121k | RETURN_IF_ERROR(_open_mapping_exprs()); |
529 | 121k | { |
530 | 121k | SCOPED_TIMER(_profile.file_reader_total_timer); |
531 | 121k | SCOPED_TIMER(_profile.file_reader_open_timer); |
532 | 121k | RETURN_IF_ERROR(_data_reader.reader->open(file_request)); |
533 | 121k | } |
534 | 121k | _file_scan_request = std::move(file_request); |
535 | 121k | RETURN_IF_ERROR(_init_reader_condition_cache(*_file_scan_request)); |
536 | 121k | return Status::OK(); |
537 | 121k | } |
538 | | |
539 | | Status _build_table_filters_from_conjuncts(); |
540 | | Status _evaluate_partition_prune_conjuncts(const VExprContextSPtrs& conjuncts, |
541 | | bool* can_filter_all); |
542 | | static bool _is_safe_to_pre_execute(const VExprContextSPtr& conjunct); |
543 | | Status _build_partition_prune_block(Block* block) const; |
544 | | Status _open_local_filter_exprs(const FileScanRequest& file_request); |
545 | | Status _init_reader_condition_cache(const FileScanRequest& file_request); |
546 | | void _finalize_reader_condition_cache(); |
547 | | bool _should_enable_condition_cache(const FileScanRequest& file_request) const; |
548 | | |
549 | 121k | Status _evaluate_constant_filters(bool* can_filter_all) { |
550 | 121k | DORIS_CHECK(can_filter_all != nullptr); |
551 | 121k | DORIS_CHECK_LE(_constant_pruning_safe_filter_count, _table_filters.size()); |
552 | 121k | *can_filter_all = false; |
553 | | // The bound was derived from the original `_conjuncts` order, which includes slotless |
554 | | // expressions omitted from `_table_filters`. Iterating only this prefix therefore cannot |
555 | | // skip an unsafe row-level predicate and pre-execute a later constant predicate. |
556 | 188k | for (size_t i = 0; i < _constant_pruning_safe_filter_count; ++i) { |
557 | 67.5k | const auto& table_filter = _table_filters[i]; |
558 | 67.5k | if (table_filter.conjunct == nullptr) { |
559 | 0 | continue; |
560 | 0 | } |
561 | 67.5k | DORIS_CHECK(_is_safe_to_pre_execute(table_filter.conjunct)); |
562 | | // RuntimeFilterExpr does not implement execute_column_impl(); it is evaluated by the |
563 | | // row-level filter path through execute_filter(). Constant split pruning uses |
564 | | // VExprContext::execute() on a one-row synthetic block, so runtime filters must not be |
565 | | // pre-executed here even when their referenced slot maps to a constant value. |
566 | 67.5k | if (table_filter.conjunct->root()->is_rf_wrapper() || |
567 | 67.5k | !_table_filter_has_only_constant_entries(table_filter)) { |
568 | 63.8k | continue; |
569 | 63.8k | } |
570 | 3.75k | Block eval_block; |
571 | 3.75k | RETURN_IF_ERROR(_build_constant_filter_block(table_filter, &eval_block)); |
572 | 3.75k | RowDescriptor row_desc; |
573 | 3.75k | RETURN_IF_ERROR(table_filter.conjunct->prepare(_runtime_state, row_desc)); |
574 | 3.75k | RETURN_IF_ERROR(table_filter.conjunct->open(_runtime_state)); |
575 | 3.75k | int result_column_id = -1; |
576 | 3.75k | RETURN_IF_ERROR(table_filter.conjunct->execute(&eval_block, &result_column_id)); |
577 | 3.75k | DORIS_CHECK(result_column_id >= 0); |
578 | 3.75k | if (_filter_result_filters_all(eval_block.get_by_position(result_column_id).column)) { |
579 | 627 | *can_filter_all = true; |
580 | 627 | return Status::OK(); |
581 | 627 | } |
582 | 3.75k | } |
583 | 120k | return Status::OK(); |
584 | 121k | } |
585 | | |
586 | 61.3k | bool _table_filter_has_only_constant_entries(const TableFilter& table_filter) const { |
587 | 61.3k | const auto& filter_entries = _data_reader.column_mapper->filter_entries(); |
588 | 61.4k | for (const auto global_index : table_filter.global_indices) { |
589 | 61.4k | const auto entry_it = filter_entries.find(global_index); |
590 | 61.6k | if (entry_it == filter_entries.end() || !entry_it->second.is_constant()) { |
591 | 57.7k | return false; |
592 | 57.7k | } |
593 | 61.4k | } |
594 | 3.60k | return !table_filter.global_indices.empty(); |
595 | 61.3k | } |
596 | | |
597 | 3.80k | Status _build_constant_filter_block(const TableFilter& table_filter, Block* eval_block) { |
598 | 3.80k | DORIS_CHECK(eval_block != nullptr); |
599 | 3.80k | eval_block->clear(); |
600 | 3.80k | const auto& mappings = _data_reader.column_mapper->mappings(); |
601 | 3.80k | const auto& filter_entries = _data_reader.column_mapper->filter_entries(); |
602 | 3.80k | DORIS_CHECK(mappings.size() == _projected_columns.size()); |
603 | 15.0k | for (size_t column_idx = 0; column_idx < mappings.size(); ++column_idx) { |
604 | 11.2k | const auto global_index = GlobalIndex(column_idx); |
605 | 11.2k | const auto& mapping = mappings[column_idx]; |
606 | 11.2k | const auto entry_it = filter_entries.find(global_index); |
607 | 11.2k | const bool referenced_by_filter = |
608 | 11.2k | std::find(table_filter.global_indices.begin(), |
609 | 11.2k | table_filter.global_indices.end(), |
610 | 11.2k | global_index) != table_filter.global_indices.end(); |
611 | 11.2k | if (referenced_by_filter && entry_it != filter_entries.end() && |
612 | 11.2k | entry_it->second.is_constant()) { |
613 | 3.85k | ColumnPtr constant_column; |
614 | 3.85k | RETURN_IF_ERROR(_materialize_constant_filter_column( |
615 | 3.85k | entry_it->second.constant_index(), &constant_column)); |
616 | 3.85k | eval_block->insert({std::move(constant_column), mapping.table_type, |
617 | 3.85k | mapping.table_column_name}); |
618 | 7.39k | } else { |
619 | 7.39k | eval_block->insert({mapping.table_type->create_column_const_with_default_value(1), |
620 | 7.39k | mapping.table_type, mapping.table_column_name}); |
621 | 7.39k | } |
622 | 11.2k | } |
623 | 3.80k | return Status::OK(); |
624 | 3.80k | } |
625 | | |
626 | 3.85k | Status _materialize_constant_filter_column(ConstantIndex constant_index, ColumnPtr* column) { |
627 | 3.85k | DORIS_CHECK(column != nullptr); |
628 | 3.85k | const auto& constant_entry = _data_reader.column_mapper->constant_map().get(constant_index); |
629 | 3.85k | DORIS_CHECK(constant_entry.expr != nullptr); |
630 | 3.85k | DORIS_CHECK(constant_entry.type != nullptr); |
631 | 3.85k | RowDescriptor row_desc; |
632 | 3.85k | RETURN_IF_ERROR(constant_entry.expr->prepare(_runtime_state, row_desc)); |
633 | 3.85k | RETURN_IF_ERROR(constant_entry.expr->open(_runtime_state)); |
634 | 3.85k | Block eval_block; |
635 | 3.85k | eval_block.insert({constant_entry.type->create_column_const_with_default_value(1), |
636 | 3.85k | constant_entry.type, "__table_reader_constant_filter"}); |
637 | 3.85k | int result_column_id = -1; |
638 | 3.85k | RETURN_IF_ERROR(constant_entry.expr->execute(&eval_block, &result_column_id)); |
639 | 3.85k | DORIS_CHECK(result_column_id >= 0); |
640 | 3.85k | *column = eval_block.get_by_position(result_column_id).column; |
641 | 3.85k | DORIS_CHECK((*column)->size() == 1); |
642 | 3.85k | return Status::OK(); |
643 | 3.85k | } |
644 | | |
645 | 3.80k | static bool _filter_result_filters_all(const ColumnPtr& filter_column) { |
646 | 3.80k | DORIS_CHECK(filter_column.get() != nullptr); |
647 | 3.80k | DORIS_CHECK(filter_column->size() == 1); |
648 | 3.80k | return !filter_column->get_bool(0); |
649 | 3.80k | } |
650 | | |
651 | 120k | virtual Status customize_file_scan_request(FileScanRequest* file_request) { |
652 | 120k | return _append_delete_predicate(file_request); |
653 | 120k | } |
654 | | |
655 | 424k | bool _is_table_level_count_active() const { return _remaining_table_level_count >= 0; } |
656 | | |
657 | 242k | bool _is_file_level_count_active() const { return _remaining_file_level_count >= 0; } |
658 | | |
659 | 3.16k | Status _materialize_count_rows(size_t rows, Block* block) const { |
660 | 3.16k | DORIS_CHECK(block != nullptr); |
661 | 3.16k | DORIS_CHECK(block->columns() > 0 || rows == 0); |
662 | 6.33k | for (size_t column_idx = 0; column_idx < block->columns(); ++column_idx) { |
663 | 3.17k | auto column = block->get_by_position(column_idx).type->create_column(); |
664 | 3.17k | if (auto* nullable = check_and_get_column<ColumnNullable>(*column)) { |
665 | | // Metadata COUNT emits synthetic input rows for the unchanged upper aggregate. |
666 | | // They must be non-NULL for COUNT(nullable_col), and constructing them explicitly |
667 | | // also keeps every nullable null map boolean-valid in debug/ASAN block checks. |
668 | 3.17k | nullable->get_nested_column().insert_many_defaults(rows); |
669 | 3.17k | nullable->get_null_map_data().resize_fill(rows, 0); |
670 | 18.4E | } else { |
671 | 18.4E | column->insert_many_defaults(rows); |
672 | 18.4E | } |
673 | 3.17k | block->replace_by_position(column_idx, std::move(column)); |
674 | 3.17k | } |
675 | 3.16k | return Status::OK(); |
676 | 3.16k | } |
677 | | |
678 | 3.17k | Status _materialize_next_count_batch(int64_t* remaining_rows, Block* block) const { |
679 | 3.17k | DORIS_CHECK(remaining_rows != nullptr); |
680 | 3.17k | DORIS_CHECK(*remaining_rows > 0); |
681 | 3.17k | const int64_t batch_size = _runtime_state == nullptr |
682 | 3.17k | ? *remaining_rows |
683 | 3.17k | : static_cast<int64_t>(_runtime_state->batch_size()); |
684 | 3.17k | const auto rows = std::min(*remaining_rows, batch_size); |
685 | 3.17k | RETURN_IF_ERROR(_materialize_count_rows(cast_set<size_t>(rows), block)); |
686 | 3.17k | *remaining_rows -= rows; |
687 | 3.17k | return Status::OK(); |
688 | 3.17k | } |
689 | | |
690 | 3.39k | Status _read_count_batch(int64_t* remaining_rows, Block* block, bool* eos) { |
691 | 3.39k | DORIS_CHECK(block != nullptr); |
692 | 3.39k | DORIS_CHECK(eos != nullptr); |
693 | 3.39k | DORIS_CHECK(_push_down_agg_type == TPushAggOp::type::COUNT); |
694 | 3.39k | DORIS_CHECK(remaining_rows != nullptr); |
695 | 3.39k | DORIS_CHECK(*remaining_rows >= 0); |
696 | 3.39k | if (*remaining_rows == 0) { |
697 | 1.69k | *remaining_rows = -1; |
698 | 1.69k | _current_task.reset(); |
699 | 1.69k | *eos = true; |
700 | 1.69k | return Status::OK(); |
701 | 1.69k | } |
702 | 1.69k | RETURN_IF_ERROR(_materialize_next_count_batch(remaining_rows, block)); |
703 | 1.69k | *eos = false; |
704 | 1.69k | return Status::OK(); |
705 | 1.69k | } |
706 | | |
707 | 458 | Status _read_table_level_count(Block* block, bool* eos) { |
708 | 458 | return _read_count_batch(&_remaining_table_level_count, block, eos); |
709 | 458 | } |
710 | | |
711 | 2.93k | Status _read_file_level_count(Block* block, bool* eos) { |
712 | 2.93k | return _read_count_batch(&_remaining_file_level_count, block, eos); |
713 | 2.93k | } |
714 | | |
715 | | void _append_file_scan_column(FileScanRequest* request, LocalColumnId column_id, |
716 | 34.9k | std::vector<LocalColumnIndex>* scan_columns) { |
717 | 34.9k | DORIS_CHECK(request != nullptr); |
718 | 34.9k | DORIS_CHECK(scan_columns != nullptr); |
719 | 34.9k | FileScanRequestBuilder builder(request); |
720 | 34.9k | Status status; |
721 | 34.9k | if (scan_columns == &request->predicate_columns) { |
722 | 32.6k | status = builder.add_predicate_column(column_id); |
723 | 32.6k | } else { |
724 | 2.25k | DORIS_CHECK(scan_columns == &request->non_predicate_columns); |
725 | 2.25k | status = builder.add_non_predicate_column(column_id); |
726 | 2.25k | } |
727 | 34.9k | DORIS_CHECK(status.ok()) << status.to_string(); |
728 | 34.9k | if (column_id == LocalColumnId(ROW_POSITION_COLUMN_ID) && |
729 | 34.9k | _find_column_definition(_data_reader.file_schema, column_id) == nullptr) { |
730 | 29.4k | _data_reader.file_schema.push_back(row_position_column_definition()); |
731 | 29.4k | } |
732 | 34.9k | } |
733 | | |
734 | | // Append DeletePredicate to file scan request if there are deletes. The predicate will be evaluated in file reader level and filter out deleted rows before returning data to table reader. |
735 | 120k | Status _append_delete_predicate(FileScanRequest* request) { |
736 | 120k | DORIS_CHECK(request != nullptr); |
737 | 120k | if ((_delete_rows == nullptr || _delete_rows->empty()) && |
738 | 120k | (_deletion_vector == nullptr || _deletion_vector->isEmpty())) { |
739 | 93.1k | return Status::OK(); |
740 | 93.1k | } |
741 | 27.3k | const auto row_position_column_id = LocalColumnId(ROW_POSITION_COLUMN_ID); |
742 | 27.3k | _append_file_scan_column(request, row_position_column_id, &request->predicate_columns); |
743 | | |
744 | 27.3k | const auto block_position = request->local_positions.at(row_position_column_id); |
745 | 27.6k | auto append_predicate = [&](auto& deleted_rows) { |
746 | 27.6k | auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows); |
747 | 27.6k | delete_predicate->add_child(VSlotRef::create_shared( |
748 | 27.6k | cast_set<int>(block_position.value()), cast_set<int>(block_position.value()), |
749 | 27.6k | -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME)); |
750 | 27.6k | request->delete_conjuncts.push_back( |
751 | 27.6k | VExprContext::create_shared(std::move(delete_predicate))); |
752 | 27.6k | }; _ZZN5doris6format11TableReader24_append_delete_predicateEPNS0_15FileScanRequestEENKUlRT_E_clISt6vectorIlSaIlEEEEDaS5_ Line | Count | Source | 745 | 1.81k | auto append_predicate = [&](auto& deleted_rows) { | 746 | 1.81k | auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows); | 747 | 1.81k | delete_predicate->add_child(VSlotRef::create_shared( | 748 | 1.81k | cast_set<int>(block_position.value()), cast_set<int>(block_position.value()), | 749 | 1.81k | -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME)); | 750 | 1.81k | request->delete_conjuncts.push_back( | 751 | 1.81k | VExprContext::create_shared(std::move(delete_predicate))); | 752 | 1.81k | }; |
_ZZN5doris6format11TableReader24_append_delete_predicateEPNS0_15FileScanRequestEENKUlRT_E_clIN7roaring12Roaring64MapEEEDaS5_ Line | Count | Source | 745 | 25.8k | auto append_predicate = [&](auto& deleted_rows) { | 746 | 25.8k | auto delete_predicate = std::make_shared<DeletePredicate>(deleted_rows); | 747 | 25.8k | delete_predicate->add_child(VSlotRef::create_shared( | 748 | 25.8k | cast_set<int>(block_position.value()), cast_set<int>(block_position.value()), | 749 | 25.8k | -1, std::make_shared<DataTypeInt64>(), ROW_POSITION_COLUMN_NAME)); | 750 | 25.8k | request->delete_conjuncts.push_back( | 751 | 25.8k | VExprContext::create_shared(std::move(delete_predicate))); | 752 | 25.8k | }; |
|
753 | 27.3k | if (_delete_rows != nullptr && !_delete_rows->empty()) { |
754 | 1.81k | append_predicate(*_delete_rows); |
755 | 1.81k | } |
756 | 27.3k | if (_deletion_vector != nullptr && !_deletion_vector->isEmpty()) { |
757 | 25.8k | append_predicate(*_deletion_vector); |
758 | 25.8k | } |
759 | 27.3k | return Status::OK(); |
760 | 120k | } |
761 | | |
762 | | // Close the current concrete reader. This hook is called by both create_next_reader() and |
763 | | // close(), so it should remain idempotent. |
764 | 121k | virtual Status close_current_reader() { |
765 | 121k | _finalize_reader_condition_cache(); |
766 | 121k | { |
767 | 121k | SCOPED_TIMER(_profile.file_reader_total_timer); |
768 | 121k | SCOPED_TIMER(_profile.file_reader_close_timer); |
769 | 121k | RETURN_IF_ERROR(_data_reader.reader->close()); |
770 | 121k | } |
771 | 121k | _data_reader.reader.reset(); |
772 | 121k | if (_data_reader.column_mapper != nullptr) { |
773 | 121k | _data_reader.column_mapper->clear(); |
774 | 121k | _data_reader.column_mapper.reset(); |
775 | 121k | } |
776 | 121k | _table_filters.clear(); |
777 | 121k | _constant_pruning_safe_filter_count = 0; |
778 | 121k | _data_reader.file_schema.clear(); |
779 | 121k | _data_reader.file_block_layout.clear(); |
780 | 121k | _data_reader.block_template.clear(); |
781 | 121k | _file_scan_request.reset(); |
782 | 121k | _current_task.reset(); |
783 | 121k | _current_file_description.reset(); |
784 | 121k | _current_reader_reached_eof = false; |
785 | 121k | return Status::OK(); |
786 | 121k | } |
787 | | |
788 | 7.12k | void _record_scan_rows(size_t rows) { |
789 | 7.12k | if (_io_ctx != nullptr && _io_ctx->file_reader_stats != nullptr) { |
790 | 7.12k | _io_ctx->file_reader_stats->read_rows += rows; |
791 | 7.12k | } |
792 | 7.12k | } |
793 | | |
794 | | // Finalize file-local block to table/global schema block. |
795 | 109k | Status finalize_chunk(Block* block, const size_t rows) { |
796 | 109k | SCOPED_TIMER(_profile.finalize_timer); |
797 | 109k | size_t idx = 0; |
798 | 109k | const auto& mappings = _data_reader.column_mapper->mappings(); |
799 | 586k | for (const auto& mapping : mappings) { |
800 | 586k | ColumnPtr column; |
801 | 586k | RETURN_IF_ERROR(_materialize_mapping_column(mapping, &_data_reader.block_template, rows, |
802 | 586k | &column, idx + 1 == mappings.size())); |
803 | 586k | block->replace_by_position(idx, IColumn::mutate(std::move(column))); |
804 | 586k | idx++; |
805 | 586k | } |
806 | 109k | RETURN_IF_ERROR(materialize_virtual_columns(block)); |
807 | | // Enforce CHAR/VARCHAR length declared by the table schema after all file-to-table |
808 | | // materialization has finished. |
809 | 109k | RETURN_IF_ERROR(_truncate_char_or_varchar_columns(block)); |
810 | 109k | return Status::OK(); |
811 | 109k | } |
812 | | |
813 | | // Materialize virtual columns in the table block, such as Iceberg _row_id and |
814 | | // _last_updated_sequence_number. This runs after normal column materialization so finalize |
815 | | // expressions can reference those virtual columns. |
816 | 78.8k | virtual Status materialize_virtual_columns(Block* table_block) { return Status::OK(); } |
817 | | |
818 | | #ifndef NDEBUG |
819 | 109k | Status _check_file_block_columns(std::string_view stage, size_t rows) { |
820 | 109k | DORIS_CHECK(_data_reader.block_template.columns() == _data_reader.file_block_layout.size()); |
821 | 679k | for (size_t idx = 0; idx < _data_reader.block_template.columns(); ++idx) { |
822 | 570k | const auto& file_block_column = _data_reader.file_block_layout[idx]; |
823 | 570k | const auto& column_with_type = _data_reader.block_template.get_by_position(idx); |
824 | 570k | const auto* column = column_with_type.column.get(); |
825 | 570k | try { |
826 | 570k | if (column == nullptr) { |
827 | 0 | auto st = Status::InternalError( |
828 | 0 | "Invalid file block column {} at {}: file_column_id={}, name='{}', " |
829 | 0 | "type={}, column=null, expected_rows={}, reader={}", |
830 | 0 | idx, stage, file_block_column.file_column_id.value(), |
831 | 0 | file_block_column.name, |
832 | 0 | file_block_column.type == nullptr ? "null" |
833 | 0 | : file_block_column.type->get_name(), |
834 | 0 | rows, debug_string()); |
835 | 0 | LOG(WARNING) << st; |
836 | 0 | return st; |
837 | 0 | } |
838 | 570k | column->sanity_check(); |
839 | 570k | auto st = column_with_type.check_type_and_column_match(); |
840 | 570k | if (!st.ok()) { |
841 | 0 | auto contextual_status = Status::InternalError( |
842 | 0 | "Invalid file block column {} at {}: file_column_id={}, name='{}', " |
843 | 0 | "type={}, column={}, column_size={}, expected_rows={}, error={}, " |
844 | 0 | "reader={}", |
845 | 0 | idx, stage, file_block_column.file_column_id.value(), |
846 | 0 | file_block_column.name, |
847 | 0 | file_block_column.type == nullptr ? "null" |
848 | 0 | : file_block_column.type->get_name(), |
849 | 0 | column->get_name(), column->size(), rows, st.to_string(), |
850 | 0 | debug_string()); |
851 | 0 | LOG(WARNING) << contextual_status; |
852 | 0 | return contextual_status; |
853 | 0 | } |
854 | 570k | } catch (const Exception& e) { |
855 | 0 | auto st = Status::InternalError( |
856 | 0 | "Invalid file block column {} at {}: file_column_id={}, name='{}', " |
857 | 0 | "type={}, column={}, column_size={}, expected_rows={}, error={}, " |
858 | 0 | "reader={}", |
859 | 0 | idx, stage, file_block_column.file_column_id.value(), |
860 | 0 | file_block_column.name, |
861 | 0 | file_block_column.type == nullptr ? "null" |
862 | 0 | : file_block_column.type->get_name(), |
863 | 0 | column == nullptr ? "null" : column->get_name(), |
864 | 0 | column == nullptr ? 0 : column->size(), rows, e.to_string(), |
865 | 0 | debug_string()); |
866 | 0 | LOG(WARNING) << st; |
867 | 0 | return st; |
868 | 0 | } catch (const std::exception& e) { |
869 | 0 | auto st = Status::InternalError( |
870 | 0 | "Invalid file block column {} at {}: file_column_id={}, name='{}', " |
871 | 0 | "type={}, column={}, column_size={}, expected_rows={}, error={}, " |
872 | 0 | "reader={}", |
873 | 0 | idx, stage, file_block_column.file_column_id.value(), |
874 | 0 | file_block_column.name, |
875 | 0 | file_block_column.type == nullptr ? "null" |
876 | 0 | : file_block_column.type->get_name(), |
877 | 0 | column == nullptr ? "null" : column->get_name(), |
878 | 0 | column == nullptr ? 0 : column->size(), rows, e.what(), debug_string()); |
879 | 0 | LOG(WARNING) << st; |
880 | 0 | return st; |
881 | 0 | } |
882 | 570k | } |
883 | 109k | return Status::OK(); |
884 | 109k | } |
885 | | |
886 | 109k | Status _check_table_block_columns(std::string_view stage, const Block* block, size_t rows) { |
887 | 109k | DORIS_CHECK(block != nullptr); |
888 | 109k | DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size()); |
889 | 695k | for (size_t idx = 0; idx < block->columns(); ++idx) { |
890 | 586k | const auto& mapping = _data_reader.column_mapper->mappings()[idx]; |
891 | 586k | const auto& column_with_type = block->get_by_position(idx); |
892 | 586k | const auto* column = column_with_type.column.get(); |
893 | 586k | try { |
894 | 586k | if (column == nullptr) { |
895 | 0 | auto st = Status::InternalError( |
896 | 0 | "Invalid table block column {} at {}: table_column='{}', " |
897 | 0 | "global_index={}, type={}, column=null, expected_rows={}, mapping={}", |
898 | 0 | idx, stage, mapping.table_column_name, mapping.global_index.value(), |
899 | 0 | mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(), |
900 | 0 | rows, mapping.debug_string()); |
901 | 0 | LOG(WARNING) << st; |
902 | 0 | return st; |
903 | 0 | } |
904 | 586k | column->sanity_check(); |
905 | 586k | auto st = column_with_type.check_type_and_column_match(); |
906 | 586k | if (!st.ok()) { |
907 | 0 | auto contextual_status = Status::InternalError( |
908 | 0 | "Invalid table block column {} at {}: table_column='{}', " |
909 | 0 | "global_index={}, type={}, column={}, column_size={}, " |
910 | 0 | "expected_rows={}, error={}, mapping={}", |
911 | 0 | idx, stage, mapping.table_column_name, mapping.global_index.value(), |
912 | 0 | mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(), |
913 | 0 | column->get_name(), column->size(), rows, st.to_string(), |
914 | 0 | mapping.debug_string()); |
915 | 0 | LOG(WARNING) << contextual_status; |
916 | 0 | return contextual_status; |
917 | 0 | } |
918 | 586k | } catch (const Exception& e) { |
919 | 0 | auto st = Status::InternalError( |
920 | 0 | "Invalid table block column {} at {}: table_column='{}', global_index={}, " |
921 | 0 | "type={}, column={}, column_size={}, expected_rows={}, error={}, " |
922 | 0 | "mapping={}", |
923 | 0 | idx, stage, mapping.table_column_name, mapping.global_index.value(), |
924 | 0 | mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(), |
925 | 0 | column == nullptr ? "null" : column->get_name(), |
926 | 0 | column == nullptr ? 0 : column->size(), rows, e.to_string(), |
927 | 0 | mapping.debug_string()); |
928 | 0 | LOG(WARNING) << st; |
929 | 0 | return st; |
930 | 0 | } catch (const std::exception& e) { |
931 | 0 | auto st = Status::InternalError( |
932 | 0 | "Invalid table block column {} at {}: table_column='{}', global_index={}, " |
933 | 0 | "type={}, column={}, column_size={}, expected_rows={}, error={}, " |
934 | 0 | "mapping={}", |
935 | 0 | idx, stage, mapping.table_column_name, mapping.global_index.value(), |
936 | 0 | mapping.table_type == nullptr ? "null" : mapping.table_type->get_name(), |
937 | 0 | column == nullptr ? "null" : column->get_name(), |
938 | 0 | column == nullptr ? 0 : column->size(), rows, e.what(), |
939 | 0 | mapping.debug_string()); |
940 | 0 | LOG(WARNING) << st; |
941 | 0 | return st; |
942 | 0 | } |
943 | 586k | } |
944 | 109k | return Status::OK(); |
945 | 109k | } |
946 | | #endif |
947 | | |
948 | 109k | Status _truncate_char_or_varchar_columns(Block* block) { |
949 | 109k | DORIS_CHECK(block != nullptr); |
950 | 109k | if (_runtime_state == nullptr || |
951 | 109k | !_runtime_state->query_options().truncate_char_or_varchar_columns) { |
952 | 109k | return Status::OK(); |
953 | 109k | } |
954 | 8 | DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size()); |
955 | 44 | for (size_t idx = 0; idx < _data_reader.column_mapper->mappings().size(); ++idx) { |
956 | 36 | const auto& mapping = _data_reader.column_mapper->mappings()[idx]; |
957 | 36 | if (!_should_truncate_char_or_varchar_column(mapping)) { |
958 | 12 | continue; |
959 | 12 | } |
960 | 24 | const auto target_len = |
961 | 24 | assert_cast<const DataTypeString*>(remove_nullable(mapping.table_type).get()) |
962 | 24 | ->len(); |
963 | 24 | _truncate_char_or_varchar_column(block, idx, target_len); |
964 | 24 | } |
965 | 8 | return Status::OK(); |
966 | 109k | } |
967 | | |
968 | | // Return true when the table schema has a bounded CHAR/VARCHAR length that is stricter than |
969 | | // the file-side type. Examples: |
970 | | // - table VARCHAR(10), file VARCHAR(20): truncate to 10; |
971 | | // - table VARCHAR(10), file STRING: truncate to 10 because STRING has no declared bound; |
972 | | // - table STRING, any file type: no truncation because the target has no bound. |
973 | 41 | static bool _should_truncate_char_or_varchar_column(const ColumnMapping& mapping) { |
974 | 41 | if (mapping.table_type == nullptr) { |
975 | 0 | return false; |
976 | 0 | } |
977 | 41 | const auto table_type = remove_nullable(mapping.table_type); |
978 | 41 | const auto primitive_type = table_type->get_primitive_type(); |
979 | 41 | if (primitive_type != TYPE_VARCHAR && primitive_type != TYPE_CHAR) { |
980 | 13 | return false; |
981 | 13 | } |
982 | 28 | const auto target_len = assert_cast<const DataTypeString*>(table_type.get())->len(); |
983 | 28 | if (target_len <= 0) { |
984 | 0 | return false; |
985 | 0 | } |
986 | 28 | if (mapping.file_type == nullptr) { |
987 | 0 | return true; |
988 | 0 | } |
989 | 28 | const auto file_type = remove_nullable(mapping.file_type); |
990 | 28 | DORIS_CHECK(file_type != nullptr); |
991 | 28 | int file_len = -1; |
992 | 28 | if (file_type->get_primitive_type() == TYPE_VARCHAR || |
993 | 28 | file_type->get_primitive_type() == TYPE_CHAR || |
994 | 28 | file_type->get_primitive_type() == TYPE_STRING) { |
995 | 27 | file_len = assert_cast<const DataTypeString*>(file_type.get())->len(); |
996 | 27 | } |
997 | | |
998 | 28 | return file_len < 0 || target_len < file_len; |
999 | 28 | } |
1000 | | |
1001 | | // Truncate a materialized CHAR/VARCHAR column in place by reusing the vectorized substring |
1002 | | // implementation: substring(column, 1, len). Nullable columns are unwrapped before substring |
1003 | | // execution and wrapped back with the original null map afterward, because substring operates |
1004 | | // on the nested string payload only. |
1005 | 25 | static void _truncate_char_or_varchar_column(Block* block, size_t idx, int len) { |
1006 | 25 | DORIS_CHECK(block != nullptr); |
1007 | 25 | auto int_type = std::make_shared<DataTypeInt32>(); |
1008 | 25 | const auto num_columns_without_result = cast_set<uint32_t>(block->columns()); |
1009 | 25 | auto& target = block->get_by_position(idx); |
1010 | 25 | const bool is_nullable = target.type->is_nullable(); |
1011 | 25 | ColumnPtr input_column = target.column; |
1012 | 25 | ColumnPtr null_map_column; |
1013 | 25 | if (is_nullable) { |
1014 | 25 | const auto* nullable_column = assert_cast<const ColumnNullable*>(target.column.get()); |
1015 | 25 | input_column = nullable_column->get_nested_column_ptr(); |
1016 | 25 | null_map_column = nullable_column->get_null_map_column_ptr(); |
1017 | 25 | } |
1018 | 25 | block->replace_by_position(idx, std::move(input_column)); |
1019 | 25 | block->insert({int_type->create_column_const(block->rows(), to_field<TYPE_INT>(1)), |
1020 | 25 | int_type, "const 1"}); |
1021 | 25 | block->insert({int_type->create_column_const(block->rows(), to_field<TYPE_INT>(len)), |
1022 | 25 | int_type, "const len"}); |
1023 | 25 | block->insert({nullptr, std::make_shared<DataTypeString>(), "result"}); |
1024 | | |
1025 | 25 | ColumnNumbers temp_arguments(3); |
1026 | 25 | temp_arguments[0] = cast_set<uint32_t>(idx); |
1027 | 25 | temp_arguments[1] = num_columns_without_result; |
1028 | 25 | temp_arguments[2] = num_columns_without_result + 1; |
1029 | 25 | const uint32_t result_column_id = num_columns_without_result + 2; |
1030 | 25 | SubstringUtil::substring_execute(*block, temp_arguments, result_column_id, block->rows()); |
1031 | | |
1032 | 25 | ColumnPtr result_column = block->get_by_position(result_column_id).column; |
1033 | 25 | if (is_nullable) { |
1034 | 25 | result_column = ColumnNullable::create(std::move(result_column), null_map_column); |
1035 | 25 | } |
1036 | 25 | block->replace_by_position(idx, std::move(result_column)); |
1037 | 25 | block->erase_tail(num_columns_without_result); |
1038 | 25 | } |
1039 | | |
1040 | 120k | Status _try_materialize_aggregate_pushdown_rows(Block* block, bool* pushed_down) { |
1041 | 120k | DORIS_CHECK(block != nullptr); |
1042 | 120k | DORIS_CHECK(pushed_down != nullptr); |
1043 | 120k | *pushed_down = false; |
1044 | 120k | block->clear_column_data(_projected_columns.size()); |
1045 | 120k | _aggregate_pushdown_tried = true; |
1046 | 120k | if (!_supports_aggregate_pushdown(_push_down_agg_type)) { |
1047 | 119k | return Status::OK(); |
1048 | 119k | } |
1049 | | |
1050 | 1.40k | FileAggregateRequest file_request; |
1051 | 1.40k | RETURN_IF_ERROR(_build_file_aggregate_request(_push_down_agg_type, &file_request)); |
1052 | 1.40k | FileAggregateResult file_result; |
1053 | 1.40k | Status status; |
1054 | 1.40k | { |
1055 | 1.40k | SCOPED_TIMER(_profile.file_reader_total_timer); |
1056 | 1.40k | SCOPED_TIMER(_profile.file_reader_aggregate_timer); |
1057 | 1.40k | status = _data_reader.reader->get_aggregate_result(file_request, &file_result); |
1058 | 1.40k | } |
1059 | 1.40k | if (status.is<ErrorCode::NOT_IMPLEMENTED_ERROR>()) { |
1060 | 13 | return Status::OK(); |
1061 | 13 | } |
1062 | 1.39k | RETURN_IF_ERROR(status); |
1063 | 1.47k | if (_push_down_agg_type == TPushAggOp::type::COUNT) { |
1064 | 1.47k | DORIS_CHECK(file_result.count >= 0); |
1065 | | // The upper aggregate consumes synthetic input rows, but emitting the whole metadata |
1066 | | // count in one block bypasses the runtime batch contract and can allocate by file size. |
1067 | | // Keep the remaining cardinality as split state and expose at most one batch per call. |
1068 | 1.47k | _remaining_file_level_count = file_result.count; |
1069 | 1.47k | _current_split_uses_metadata_count = true; |
1070 | 1.47k | if (_remaining_file_level_count > 0) { |
1071 | 1.47k | RETURN_IF_ERROR(_materialize_next_count_batch(&_remaining_file_level_count, block)); |
1072 | 1.47k | } |
1073 | 18.4E | } else { |
1074 | 18.4E | RETURN_IF_ERROR( |
1075 | 18.4E | _materialize_aggregate_pushdown_rows(_push_down_agg_type, file_result, block)); |
1076 | 18.4E | } |
1077 | 1.38k | *pushed_down = true; |
1078 | 1.38k | RETURN_IF_ERROR(close_current_reader()); |
1079 | 1.38k | return Status::OK(); |
1080 | 1.38k | } |
1081 | | |
1082 | 122k | virtual bool _supports_aggregate_pushdown(TPushAggOp::type agg_type) const { |
1083 | | // Only COUNT and MIN/MAX can be push down. |
1084 | 122k | if (agg_type != TPushAggOp::type::COUNT && agg_type != TPushAggOp::type::MINMAX) { |
1085 | 118k | return false; |
1086 | 118k | } |
1087 | | // Aggregate pushdown returns reduced synthetic rows and may close the physical reader |
1088 | | // before the next scheduler turn. If a runtime filter is still pending, those rows could |
1089 | | // escape before the filter arrives and cannot later be reconstructed from real file rows. |
1090 | | // This is the same irreversibility constraint as table-level metadata COUNT, and applies |
1091 | | // to COUNT and MIN/MAX for Parquet/ORC as well as COUNT for text readers. |
1092 | 4.10k | if (!_all_runtime_filters_applied_for_split) { |
1093 | 3 | return false; |
1094 | 3 | } |
1095 | | // Scanner owns the original conjunct list and evaluates it after TableReader finalizes |
1096 | | // rows. Even a slotless conjunct that cannot become a TableFilter must see every source |
1097 | | // row before an aggregate reduces the stream to synthetic COUNT/MINMAX rows. |
1098 | 4.10k | if (!_conjuncts.empty()) { |
1099 | 5 | return false; |
1100 | 5 | } |
1101 | | // Only support aggregate pushdown when there is no delete or filter, so |
1102 | | // the reduced rows consumed by the upper aggregate remain semantically equivalent to a |
1103 | | // normal scan. |
1104 | 4.09k | if ((_delete_rows != nullptr && !_delete_rows->empty()) || |
1105 | 4.49k | (_deletion_vector != nullptr && !_deletion_vector->isEmpty())) { |
1106 | 1.09k | return false; |
1107 | 1.09k | } |
1108 | 3.00k | if (!_table_filters.empty()) { |
1109 | 0 | return false; |
1110 | 0 | } |
1111 | 3.42k | if (agg_type == TPushAggOp::type::COUNT) { |
1112 | | // Old FEs do not serialize push_down_count_slot_ids. During the supported BE-first |
1113 | | // rolling upgrade, nullopt therefore means "COUNT semantics are unknown", not |
1114 | | // COUNT(*). Fall back to reading rows until the FE explicitly sends either an empty |
1115 | | // list for COUNT(*) or one slot for COUNT(col). |
1116 | 3.42k | if (!_push_down_count_columns.has_value()) { |
1117 | 3 | return false; |
1118 | 3 | } |
1119 | | // COUNT(*) needs no column metadata. COUNT(col) currently supports one direct file |
1120 | | // column; multiple COUNT arguments fall back to the normal scan so every upper |
1121 | | // aggregate receives the original rows. |
1122 | 3.42k | if (_push_down_count_columns->empty()) { |
1123 | 2.59k | return true; |
1124 | 2.59k | } |
1125 | 827 | if (_push_down_count_columns->size() != 1) { |
1126 | 37 | return false; |
1127 | 37 | } |
1128 | 790 | const auto& mapping = _push_down_count_mapping(); |
1129 | | // Metadata COUNT skips TableReader's normal materialization path. Only a trivial |
1130 | | // mapping is safe: for example, a nullable Parquet INT mapped to a NOT NULL table |
1131 | | // BIGINT normally needs both an INT->BIGINT cast and nullability validation. Counting |
1132 | | // footer values directly would bypass both operations and could hide invalid data. |
1133 | 790 | return mapping.file_local_id.has_value() && mapping.file_type != nullptr && |
1134 | 790 | mapping.table_type != nullptr && mapping.is_trivial && |
1135 | 790 | mapping.virtual_column_type == TableVirtualColumnType::INVALID && |
1136 | 790 | mapping.default_expr == nullptr; |
1137 | 827 | } |
1138 | | // For MIN/MAX, only support direct file-to-table column mappings. The two emitted rows |
1139 | | // must be enough for the upper MIN/MAX aggregate without evaluating default expressions or |
1140 | | // virtual columns. |
1141 | 18.4E | for (const auto& mapping : _data_reader.column_mapper->mappings()) { |
1142 | 149 | if (!mapping.file_local_id.has_value() || |
1143 | 149 | mapping.virtual_column_type != TableVirtualColumnType::INVALID || |
1144 | 149 | mapping.default_expr != nullptr || mapping.file_type == nullptr || |
1145 | 149 | mapping.table_type == nullptr) { |
1146 | 9 | return false; |
1147 | 9 | } |
1148 | 140 | if (!_can_push_down_minmax_for_mapping(mapping)) { |
1149 | 46 | return false; |
1150 | 46 | } |
1151 | 140 | } |
1152 | 18.4E | return true; |
1153 | 18.4E | } |
1154 | | |
1155 | 581k | static ColumnPtr _detach_column(ColumnPtr column) { |
1156 | 581k | DORIS_CHECK(column.get() != nullptr); |
1157 | 581k | return IColumn::mutate(std::move(column)); |
1158 | 581k | } |
1159 | | |
1160 | 86.3k | static ColumnPtr _take_and_detach_block_column(Block* block, int position) { |
1161 | 86.3k | DORIS_CHECK(block != nullptr); |
1162 | 86.3k | DORIS_CHECK(position >= 0 && position < static_cast<int>(block->columns())); |
1163 | 86.3k | auto& source = block->get_by_position(position); |
1164 | 86.3k | ColumnPtr column = source.column; |
1165 | | // The final mapping no longer needs the file block. Release its COW owner before mutate(), |
1166 | | // otherwise nested MAP/STRING columns are deep-copied and a multi-GB payload can OOM. |
1167 | 86.3k | block->replace_by_position(position, source.type->create_column()); |
1168 | 86.3k | return _detach_column(std::move(column)); |
1169 | 86.3k | } |
1170 | | |
1171 | | static Status _align_column_nullability(ColumnPtr* column, const DataTypePtr& table_type, |
1172 | 160k | const NullMap* nullable_parent_null_map = nullptr) { |
1173 | 160k | DORIS_CHECK(column != nullptr); |
1174 | 160k | DORIS_CHECK(column->get() != nullptr); |
1175 | 160k | DORIS_CHECK(table_type != nullptr); |
1176 | | // Must return non-const column |
1177 | 160k | *column = (*column)->convert_to_full_column_if_const(); |
1178 | 160k | if (table_type->is_nullable()) { |
1179 | 71.3k | const auto& nested_type = |
1180 | 71.3k | assert_cast<const DataTypeNullable&>(*table_type).get_nested_type(); |
1181 | 71.3k | if (!(*column)->is_nullable()) { |
1182 | 10 | RETURN_IF_ERROR( |
1183 | 10 | _align_column_nullability(column, nested_type, nullable_parent_null_map)); |
1184 | 10 | *column = make_nullable(*column); |
1185 | 10 | return Status::OK(); |
1186 | 10 | } |
1187 | 71.3k | const auto& nullable_column = assert_cast<const ColumnNullable&>(**column); |
1188 | 71.3k | ColumnPtr nested_column = nullable_column.get_nested_column_ptr(); |
1189 | 71.3k | NullMap combined_null_map; |
1190 | 71.3k | const NullMap* nested_parent_null_map = &nullable_column.get_null_map_data(); |
1191 | 71.3k | if (nullable_parent_null_map != nullptr) { |
1192 | 31.6k | const auto& own_null_map = nullable_column.get_null_map_data(); |
1193 | 31.6k | DORIS_CHECK(nullable_parent_null_map->size() == own_null_map.size()); |
1194 | | // Required descendants are hidden when either this nullable container or any |
1195 | | // inherited nullable ancestor masks the row, so preserve the union recursively. |
1196 | 31.6k | combined_null_map.resize(own_null_map.size()); |
1197 | 352k | for (size_t i = 0; i < own_null_map.size(); ++i) { |
1198 | 321k | combined_null_map[i] = own_null_map[i] || (*nullable_parent_null_map)[i]; |
1199 | 321k | } |
1200 | 31.6k | nested_parent_null_map = &combined_null_map; |
1201 | 31.6k | } |
1202 | 71.3k | RETURN_IF_ERROR( |
1203 | 71.3k | _align_column_nullability(&nested_column, nested_type, nested_parent_null_map)); |
1204 | 71.3k | *column = ColumnNullable::create(nested_column, |
1205 | 71.3k | nullable_column.get_null_map_column_ptr()); |
1206 | 71.3k | return Status::OK(); |
1207 | 71.3k | } |
1208 | 89.1k | if ((*column)->is_nullable()) { |
1209 | 8.62k | const auto& nullable_column = assert_cast<const ColumnNullable&>(**column); |
1210 | 8.62k | if (nullable_column.has_null()) { |
1211 | 28 | const auto& null_map = nullable_column.get_null_map_data(); |
1212 | 28 | if (nullable_parent_null_map == nullptr || |
1213 | 28 | nullable_parent_null_map->size() != null_map.size()) { |
1214 | 1 | return Status::InternalError( |
1215 | 1 | "Default expression produced NULL for non-nullable table column"); |
1216 | 1 | } |
1217 | 64 | for (size_t i = 0; i < null_map.size(); ++i) { |
1218 | | // A required child may contain a physical NULL placeholder only when its |
1219 | | // nullable parent masks that row from the logical value. |
1220 | 38 | if (null_map[i] && !(*nullable_parent_null_map)[i]) { |
1221 | 1 | return Status::InternalError( |
1222 | 1 | "Default expression produced NULL for non-nullable table column"); |
1223 | 1 | } |
1224 | 38 | } |
1225 | 27 | } |
1226 | 8.62k | ColumnPtr nested_column = nullable_column.get_nested_column_ptr(); |
1227 | 8.62k | RETURN_IF_ERROR(_align_column_nullability(&nested_column, table_type, |
1228 | 8.62k | nullable_parent_null_map)); |
1229 | 8.62k | *column = nested_column; |
1230 | 8.62k | return Status::OK(); |
1231 | 8.62k | } |
1232 | 80.4k | if (const auto* array_type = typeid_cast<const DataTypeArray*>(table_type.get())) { |
1233 | 230 | const auto& array_column = assert_cast<const ColumnArray&>(**column); |
1234 | 230 | ColumnPtr nested_column = array_column.get_data_ptr(); |
1235 | 230 | NullMap descendant_parent_null_map; |
1236 | | // Collection entries use offset coordinates, so inherited row masks must be projected |
1237 | | // only when a required descendant can consume them. This avoids scratch proportional |
1238 | | // to all array entries for the common all-required schema. |
1239 | 230 | const NullMap* descendant_parent_null_map_ptr = nullptr; |
1240 | 230 | if (_requires_collection_parent_null_map( |
1241 | 230 | nullable_parent_null_map, nested_column, array_type->get_nested_type(), |
1242 | 230 | array_column.size(), array_column.get_offsets())) { |
1243 | 1 | descendant_parent_null_map_ptr = |
1244 | 1 | _project_collection_parent_null_map_for_hidden_entries( |
1245 | 1 | nullptr, nullable_parent_null_map, array_column.size(), |
1246 | 1 | array_column.get_offsets(), nested_column->size(), |
1247 | 1 | &descendant_parent_null_map); |
1248 | 1 | } |
1249 | 230 | RETURN_IF_ERROR(_align_column_nullability(&nested_column, array_type->get_nested_type(), |
1250 | 230 | descendant_parent_null_map_ptr)); |
1251 | 230 | *column = ColumnArray::create(nested_column, array_column.get_offsets_ptr()); |
1252 | 230 | return Status::OK(); |
1253 | 230 | } |
1254 | 80.2k | if (const auto* map_type = typeid_cast<const DataTypeMap*>(table_type.get())) { |
1255 | 181 | const auto& map_column = assert_cast<const ColumnMap&>(**column); |
1256 | 181 | ColumnPtr key_column = map_column.get_keys_ptr(); |
1257 | 181 | ColumnPtr value_column = map_column.get_values_ptr(); |
1258 | 181 | NullMap descendant_parent_null_map; |
1259 | 181 | const NullMap* descendant_parent_null_map_ptr = nullptr; |
1260 | 181 | if (_requires_collection_parent_null_map(nullable_parent_null_map, key_column, |
1261 | 181 | map_type->get_key_type(), map_column.size(), |
1262 | 181 | map_column.get_offsets()) || |
1263 | 181 | _requires_collection_parent_null_map(nullable_parent_null_map, value_column, |
1264 | 181 | map_type->get_value_type(), map_column.size(), |
1265 | 181 | map_column.get_offsets())) { |
1266 | | // Keys and values share offsets, so one projected mask safely covers both streams. |
1267 | 0 | descendant_parent_null_map_ptr = |
1268 | 0 | _project_collection_parent_null_map_for_hidden_entries( |
1269 | 0 | nullptr, nullable_parent_null_map, map_column.size(), |
1270 | 0 | map_column.get_offsets(), key_column->size(), |
1271 | 0 | &descendant_parent_null_map); |
1272 | 0 | } |
1273 | 181 | RETURN_IF_ERROR(_align_column_nullability(&key_column, map_type->get_key_type(), |
1274 | 181 | descendant_parent_null_map_ptr)); |
1275 | 181 | RETURN_IF_ERROR(_align_column_nullability(&value_column, map_type->get_value_type(), |
1276 | 181 | descendant_parent_null_map_ptr)); |
1277 | 181 | *column = ColumnMap::create(key_column, value_column, map_column.get_offsets_ptr()); |
1278 | 181 | return Status::OK(); |
1279 | 181 | } |
1280 | 80.0k | if (const auto* struct_type = typeid_cast<const DataTypeStruct*>(table_type.get())) { |
1281 | 8.51k | const auto& struct_column = assert_cast<const ColumnStruct&>(**column); |
1282 | 8.51k | Columns columns = struct_column.get_columns_copy(); |
1283 | 8.51k | DORIS_CHECK(columns.size() == struct_type->get_elements().size()); |
1284 | 29.1k | for (size_t i = 0; i < columns.size(); ++i) { |
1285 | 20.6k | RETURN_IF_ERROR(_align_column_nullability(&columns[i], struct_type->get_element(i), |
1286 | 20.6k | nullable_parent_null_map)); |
1287 | 20.6k | } |
1288 | 8.51k | *column = ColumnStruct::create(columns); |
1289 | 8.51k | return Status::OK(); |
1290 | 8.51k | } |
1291 | 71.5k | return Status::OK(); |
1292 | 80.0k | } |
1293 | | |
1294 | | static Status _execute_default_expr_without_root_type_check( |
1295 | | const VExprContextSPtr& default_expr, const Block* block, |
1296 | 21.6k | ColumnWithTypeAndName* result_data) { |
1297 | 21.6k | DORIS_CHECK(default_expr != nullptr); |
1298 | 21.6k | DORIS_CHECK(block != nullptr); |
1299 | 21.6k | DORIS_CHECK(result_data != nullptr); |
1300 | 21.6k | ColumnPtr result_column; |
1301 | 21.6k | Status st; |
1302 | 21.6k | RETURN_IF_CATCH_EXCEPTION({ |
1303 | 21.6k | st = default_expr->root()->execute_column_impl(default_expr.get(), block, nullptr, |
1304 | 21.6k | block->rows(), result_column); |
1305 | 21.6k | }); |
1306 | 21.6k | RETURN_IF_ERROR(st); |
1307 | 21.6k | DORIS_CHECK(result_column.get() != nullptr); |
1308 | 21.6k | if (result_column->size() != block->rows()) { |
1309 | 0 | return Status::InternalError( |
1310 | 0 | "Default expr {} return column size {} not equal to expected size {}", |
1311 | 0 | default_expr->expr_name(), result_column->size(), block->rows()); |
1312 | 0 | } |
1313 | 21.6k | result_data->column = result_column; |
1314 | 21.6k | result_data->type = default_expr->execute_type(block); |
1315 | 21.6k | result_data->name = default_expr->expr_name(); |
1316 | 21.6k | return Status::OK(); |
1317 | 21.6k | } |
1318 | | |
1319 | | Status _cast_column_to_type(ColumnPtr* column, const DataTypePtr& file_type, |
1320 | | const DataTypePtr& table_type, |
1321 | 18.2k | const std::string& column_name) const { |
1322 | 18.2k | DORIS_CHECK(column != nullptr); |
1323 | 18.2k | DORIS_CHECK(column->get() != nullptr); |
1324 | 18.2k | DORIS_CHECK(file_type != nullptr); |
1325 | 18.2k | DORIS_CHECK(table_type != nullptr); |
1326 | 18.2k | if (file_type->equals(*table_type) || |
1327 | 18.2k | remove_nullable(file_type)->equals(*remove_nullable(table_type))) { |
1328 | 8.61k | return Status::OK(); |
1329 | 8.61k | } |
1330 | | |
1331 | 9.66k | DataTypePtr input_type = file_type; |
1332 | | // Cast wrappers unwrap nullable inputs according to the declared input type, so keep the |
1333 | | // root nullability of the declared input aligned with the actual column shape. When the |
1334 | | // runtime column is nullable, also keep the cast target nullable; the caller applies the |
1335 | | // table's final nullability after value conversion. Casting a nullable runtime column |
1336 | | // directly to a non-nullable target would pass ColumnNullable to CastToImpl. |
1337 | 9.66k | if ((*column)->is_nullable() && !input_type->is_nullable()) { |
1338 | 0 | input_type = make_nullable(input_type); |
1339 | 9.66k | } else if (!(*column)->is_nullable() && input_type->is_nullable()) { |
1340 | 1 | input_type = remove_nullable(input_type); |
1341 | 1 | } |
1342 | 9.66k | DataTypePtr cast_type = table_type; |
1343 | 9.66k | if ((*column)->is_nullable() && !cast_type->is_nullable()) { |
1344 | 11 | cast_type = make_nullable(cast_type); |
1345 | 11 | } |
1346 | 9.66k | Block cast_block; |
1347 | 9.66k | cast_block.insert({*column, input_type, column_name}); |
1348 | 9.66k | auto slot_ref = VSlotRef::create_shared(0, 0, -1, input_type, column_name); |
1349 | 9.66k | auto cast_expr = Cast::create_shared(cast_type); |
1350 | 9.66k | cast_expr->add_child(std::move(slot_ref)); |
1351 | 9.66k | auto cast_ctx = VExprContext::create_shared(std::move(cast_expr)); |
1352 | 9.66k | RowDescriptor row_desc; |
1353 | 9.66k | RETURN_IF_ERROR(cast_ctx->prepare(_runtime_state, row_desc)); |
1354 | 9.66k | RETURN_IF_ERROR(cast_ctx->open(_runtime_state)); |
1355 | 9.66k | ColumnPtr cast_column; |
1356 | 9.66k | RETURN_IF_ERROR(cast_ctx->execute(&cast_block, cast_column)); |
1357 | 9.66k | *column = std::move(cast_column); |
1358 | 9.66k | return Status::OK(); |
1359 | 9.66k | } |
1360 | | |
1361 | | Status _try_materialize_scalar_cast_with_runtime_nullability(const ColumnMapping& mapping, |
1362 | | const Block* current_block, |
1363 | | ColumnPtr* column, |
1364 | 574k | bool* handled) const { |
1365 | 574k | DORIS_CHECK(column != nullptr); |
1366 | 574k | DORIS_CHECK(handled != nullptr); |
1367 | 574k | *handled = false; |
1368 | 574k | if (mapping.projection == nullptr || !mapping.file_local_id.has_value() || |
1369 | 574k | !mapping.child_mappings.empty()) { |
1370 | 188k | return Status::OK(); |
1371 | 188k | } |
1372 | | |
1373 | 386k | const auto& root = mapping.projection->root(); |
1374 | 386k | if (root == nullptr || root->node_type() != TExprNodeType::CAST_EXPR) { |
1375 | 369k | return Status::OK(); |
1376 | 369k | } |
1377 | 16.4k | DORIS_CHECK(root->get_num_children() == 1); |
1378 | 16.4k | const auto* slot = dynamic_cast<const VSlotRef*>(root->get_child(0).get()); |
1379 | 16.4k | DORIS_CHECK(slot != nullptr); |
1380 | 16.4k | DORIS_CHECK(current_block != nullptr); |
1381 | 16.4k | DORIS_CHECK(slot->column_id() >= 0); |
1382 | 16.4k | DORIS_CHECK(cast_set<size_t>(slot->column_id()) < current_block->columns()); |
1383 | 16.4k | const auto& source = current_block->get_by_position(slot->column_id()); |
1384 | 16.4k | DORIS_CHECK(source.column.get() != nullptr); |
1385 | 16.4k | DORIS_CHECK(slot->data_type() != nullptr); |
1386 | 16.4k | DORIS_CHECK(mapping.table_type != nullptr); |
1387 | 16.4k | const bool runtime_input_mismatch = |
1388 | 16.4k | source.column->is_nullable() != slot->data_type()->is_nullable(); |
1389 | 16.4k | const bool nullable_input_to_required_table = |
1390 | 16.4k | source.column->is_nullable() && !mapping.table_type->is_nullable(); |
1391 | 16.4k | if (!runtime_input_mismatch && !nullable_input_to_required_table) { |
1392 | 7.83k | return Status::OK(); |
1393 | 7.83k | } |
1394 | | |
1395 | | // File readers can return a nullable runtime column even when the physical schema marks the |
1396 | | // leaf required. A pre-built Cast binds to the declared file type and can therefore pass a |
1397 | | // ColumnNullable to a non-nullable CastToImpl. Rebuild only when that runtime shape differs |
1398 | | // from the declared input, or when a declared nullable file field maps to a required table |
1399 | | // field. Keep the cast target nullable while converting values, then let |
1400 | | // _align_column_nullability() reject an actual NULL before removing the wrapper. |
1401 | 8.57k | ColumnPtr result_column = source.column; |
1402 | 8.57k | RETURN_IF_ERROR(_cast_column_to_type(&result_column, slot->data_type(), mapping.table_type, |
1403 | 8.57k | mapping.file_column_name)); |
1404 | 8.57k | RETURN_IF_ERROR(_align_column_nullability(&result_column, mapping.table_type)); |
1405 | 8.57k | *column = _detach_column(std::move(result_column)); |
1406 | 8.57k | *handled = true; |
1407 | 8.57k | return Status::OK(); |
1408 | 8.57k | } |
1409 | | |
1410 | | Status _materialize_present_child_mapping_column( |
1411 | | const ColumnMapping& mapping, const ColumnPtr& file_column, const size_t rows, |
1412 | 21.3k | ColumnPtr* column, const NullMap* nullable_parent_null_map = nullptr) { |
1413 | 21.3k | DORIS_CHECK(column != nullptr); |
1414 | 21.3k | DORIS_CHECK(mapping.file_type != nullptr); |
1415 | 21.3k | DORIS_CHECK(mapping.table_type != nullptr); |
1416 | 21.3k | *column = file_column; |
1417 | 21.3k | if (!mapping.is_trivial) { |
1418 | 12.2k | if (!mapping.child_mappings.empty()) { |
1419 | 2.59k | RETURN_IF_ERROR(_materialize_complex_mapping_column(mapping, *column, rows, column, |
1420 | 2.59k | nullable_parent_null_map)); |
1421 | 9.69k | } else { |
1422 | 9.69k | RETURN_IF_ERROR(_cast_column_to_type(column, mapping.file_type, mapping.table_type, |
1423 | 9.69k | mapping.file_column_name)); |
1424 | 9.69k | } |
1425 | 12.2k | } |
1426 | 21.3k | RETURN_IF_ERROR( |
1427 | 21.3k | _align_column_nullability(column, mapping.table_type, nullable_parent_null_map)); |
1428 | 21.3k | return Status::OK(); |
1429 | 21.3k | } |
1430 | | |
1431 | | Status _materialize_default_or_missing_column( |
1432 | | const ColumnMapping& mapping, const Block* current_block, const size_t rows, |
1433 | 29.2k | ColumnPtr* column, const NullMap* nullable_parent_null_map = nullptr) { |
1434 | 29.2k | DORIS_CHECK(mapping.table_type != nullptr); |
1435 | 29.2k | DORIS_CHECK(column != nullptr); |
1436 | 29.2k | if (mapping.default_expr != nullptr) { |
1437 | 21.5k | Block synthetic_block; |
1438 | 21.5k | const Block* eval_block = current_block; |
1439 | 21.5k | if (eval_block == nullptr || eval_block->rows() != rows) { |
1440 | | // Nested ARRAY/MAP children use element/entry cardinality rather than the root |
1441 | | // block's row count. Iceberg initial defaults are typed literals, so a synthetic |
1442 | | // block with the desired row count is sufficient and avoids a top-level |
1443 | | // ConstantMap dependency for nested mappings. |
1444 | 2.99k | synthetic_block.insert( |
1445 | 2.99k | {mapping.table_type->create_column_const_with_default_value(rows), |
1446 | 2.99k | mapping.table_type, "__table_reader_nested_default_rows"}); |
1447 | 2.99k | eval_block = &synthetic_block; |
1448 | 2.99k | } |
1449 | 21.5k | ColumnWithTypeAndName result; |
1450 | 21.5k | RETURN_IF_ERROR(_execute_default_expr_without_root_type_check(mapping.default_expr, |
1451 | 21.5k | eval_block, &result)); |
1452 | 21.5k | ColumnPtr result_column = result.column; |
1453 | 21.5k | RETURN_IF_ERROR(_align_column_nullability(&result_column, mapping.table_type, |
1454 | 21.5k | nullable_parent_null_map)); |
1455 | 21.5k | *column = _detach_column(std::move(result_column)); |
1456 | 21.5k | return Status::OK(); |
1457 | 21.5k | } |
1458 | 7.66k | ColumnPtr result_column = mapping.table_type->create_column_const_with_default_value(rows); |
1459 | 7.66k | RETURN_IF_ERROR(_align_column_nullability(&result_column, mapping.table_type, |
1460 | 7.66k | nullable_parent_null_map)); |
1461 | 7.66k | *column = _detach_column(std::move(result_column)); |
1462 | 7.66k | return Status::OK(); |
1463 | 7.66k | } |
1464 | | |
1465 | | Status _materialize_mapping_column(const ColumnMapping& mapping, Block* current_block, |
1466 | | const size_t rows, ColumnPtr* column, |
1467 | 586k | bool take_projection_result = false) { |
1468 | 586k | if (!mapping.is_trivial && mapping.file_local_id.has_value() && |
1469 | 586k | !mapping.child_mappings.empty()) { |
1470 | 11.4k | DCHECK(mapping.projection != nullptr); |
1471 | 11.4k | int res_id; |
1472 | 11.4k | auto st = mapping.projection->execute(current_block, &res_id); |
1473 | 11.4k | if (!st.ok()) { |
1474 | 0 | return Status::InternalError( |
1475 | 0 | "Failed to execute complex mapping projection for table column '{}' " |
1476 | 0 | "(global_index={}, file_local_id={}, rows={}): {}, mapping={}", |
1477 | 0 | mapping.table_column_name, mapping.global_index.value(), |
1478 | 0 | *mapping.file_local_id, rows, st.to_string(), mapping.debug_string()); |
1479 | 0 | } |
1480 | 11.4k | ColumnPtr result_column = take_projection_result |
1481 | 11.4k | ? _take_and_detach_block_column(current_block, res_id) |
1482 | 11.4k | : current_block->get_by_position(res_id).column; |
1483 | 11.4k | RETURN_IF_ERROR( |
1484 | 11.4k | _materialize_complex_mapping_column(mapping, result_column, rows, column)); |
1485 | 11.4k | return Status::OK(); |
1486 | 11.4k | } |
1487 | 574k | bool runtime_nullability_cast_handled = false; |
1488 | 574k | RETURN_IF_ERROR(_try_materialize_scalar_cast_with_runtime_nullability( |
1489 | 574k | mapping, current_block, column, &runtime_nullability_cast_handled)); |
1490 | 574k | if (runtime_nullability_cast_handled) { |
1491 | 8.58k | return Status::OK(); |
1492 | 8.58k | } |
1493 | 566k | if (mapping.projection != nullptr) { |
1494 | 542k | int res_id; |
1495 | 542k | auto st = mapping.projection->execute(current_block, &res_id); |
1496 | 542k | if (!st.ok()) { |
1497 | 0 | std::string file_local_id = "null"; |
1498 | 0 | if (mapping.file_local_id.has_value()) { |
1499 | 0 | file_local_id = std::to_string(*mapping.file_local_id); |
1500 | 0 | } |
1501 | 0 | return Status::InternalError( |
1502 | 0 | "Failed to execute mapping projection for table column '{}' " |
1503 | 0 | "(global_index={}, file_local_id={}, rows={}): {}, mapping={}", |
1504 | 0 | mapping.table_column_name, mapping.global_index.value(), file_local_id, |
1505 | 0 | rows, st.to_string(), mapping.debug_string()); |
1506 | 0 | } |
1507 | 542k | if (take_projection_result) { |
1508 | 84.5k | *column = _take_and_detach_block_column(current_block, res_id); |
1509 | 457k | } else { |
1510 | 457k | ColumnPtr result_column = current_block->get_by_position(res_id).column; |
1511 | 457k | *column = _detach_column(std::move(result_column)); |
1512 | 457k | } |
1513 | 542k | return Status::OK(); |
1514 | 542k | } |
1515 | 24.1k | return _materialize_default_or_missing_column(mapping, current_block, rows, column); |
1516 | 566k | } |
1517 | | |
1518 | | Status _materialize_complex_mapping_column(const ColumnMapping& mapping, |
1519 | | const ColumnPtr& file_column, const size_t rows, |
1520 | | ColumnPtr* column, |
1521 | 14.0k | const NullMap* nullable_parent_null_map = nullptr) { |
1522 | 14.0k | DORIS_CHECK(mapping.table_type != nullptr); |
1523 | 14.0k | DORIS_CHECK(file_column.get() != nullptr); |
1524 | 14.0k | const auto table_type = remove_nullable(mapping.table_type); |
1525 | 14.0k | switch (table_type->get_primitive_type()) { |
1526 | 4.78k | case TYPE_STRUCT: |
1527 | 4.78k | RETURN_IF_ERROR(_materialize_struct_mapping_column(mapping, file_column, rows, column, |
1528 | 4.78k | nullable_parent_null_map)); |
1529 | 4.78k | break; |
1530 | 4.82k | case TYPE_ARRAY: |
1531 | 4.82k | RETURN_IF_ERROR(_materialize_array_mapping_column(mapping, file_column, rows, column, |
1532 | 4.82k | nullable_parent_null_map)); |
1533 | 4.82k | break; |
1534 | 4.82k | case TYPE_MAP: |
1535 | 4.39k | RETURN_IF_ERROR(_materialize_map_mapping_column(mapping, file_column, rows, column, |
1536 | 4.39k | nullable_parent_null_map)); |
1537 | 4.39k | break; |
1538 | 4.39k | default: |
1539 | 0 | *column = _detach_column(file_column); |
1540 | 0 | break; |
1541 | 14.0k | } |
1542 | 14.0k | return Status::OK(); |
1543 | 14.0k | } |
1544 | | |
1545 | | static std::vector<const ColumnMapping*> _present_child_mappings_in_file_order( |
1546 | 4.79k | const std::vector<ColumnMapping>& child_mappings) { |
1547 | 4.79k | std::vector<const ColumnMapping*> result; |
1548 | 4.79k | result.reserve(child_mappings.size()); |
1549 | 12.7k | for (const auto& child_mapping : child_mappings) { |
1550 | 12.7k | if (child_mapping.file_local_id.has_value()) { |
1551 | 7.73k | result.push_back(&child_mapping); |
1552 | 7.73k | } |
1553 | 12.7k | } |
1554 | 6.50k | std::ranges::sort(result, [](const ColumnMapping* lhs, const ColumnMapping* rhs) { |
1555 | 6.50k | DORIS_CHECK(lhs->file_local_id.has_value()); |
1556 | 6.50k | DORIS_CHECK(rhs->file_local_id.has_value()); |
1557 | 6.50k | return *lhs->file_local_id < *rhs->file_local_id; |
1558 | 6.50k | }); |
1559 | 4.79k | return result; |
1560 | 4.79k | } |
1561 | | |
1562 | | static size_t _file_child_ordinal_for_mapping( |
1563 | | const ColumnMapping& mapping, const ColumnMapping& child_mapping, |
1564 | 7.73k | const std::vector<const ColumnMapping*>& file_ordered_children) { |
1565 | 7.73k | DORIS_CHECK(child_mapping.file_local_id.has_value()); |
1566 | 7.73k | if (!mapping.projected_file_children.empty()) { |
1567 | 7.72k | const auto child_it = std::ranges::find_if( |
1568 | 15.2k | mapping.projected_file_children, [&](const ColumnDefinition& file_child) { |
1569 | 15.2k | return file_child.file_local_id() == *child_mapping.file_local_id; |
1570 | 15.2k | }); |
1571 | 7.72k | DORIS_CHECK(child_it != mapping.projected_file_children.end()); |
1572 | 7.72k | return static_cast<size_t>( |
1573 | 7.72k | std::distance(mapping.projected_file_children.begin(), child_it)); |
1574 | 7.72k | } |
1575 | 11 | const auto child_it = std::ranges::find(file_ordered_children, &child_mapping); |
1576 | 11 | DORIS_CHECK(child_it != file_ordered_children.end()); |
1577 | 11 | return static_cast<size_t>(std::distance(file_ordered_children.begin(), child_it)); |
1578 | 7.73k | } |
1579 | | |
1580 | | static std::vector<const ColumnMapping*> _child_mappings_in_table_type_order( |
1581 | 4.79k | const ColumnMapping& mapping, const DataTypeStruct& table_type) { |
1582 | 4.79k | std::vector<const ColumnMapping*> result; |
1583 | 4.79k | result.reserve(mapping.child_mappings.size()); |
1584 | 17.5k | for (size_t child_idx = 0; child_idx < table_type.get_elements().size(); ++child_idx) { |
1585 | 12.7k | const auto& child_name = table_type.get_element_name(child_idx); |
1586 | 12.7k | const auto child_it = std::ranges::find_if( |
1587 | 28.8k | mapping.child_mappings, [&](const ColumnMapping& child_mapping) { |
1588 | 28.8k | return child_mapping.table_column_name == child_name; |
1589 | 28.8k | }); |
1590 | 12.7k | DORIS_CHECK(child_it != mapping.child_mappings.end()) |
1591 | 0 | << mapping.debug_string() << ", table_child_name=" << child_name; |
1592 | 12.7k | result.push_back(&*child_it); |
1593 | 12.7k | } |
1594 | 4.79k | return result; |
1595 | 4.79k | } |
1596 | | |
1597 | | static const IColumn* _nested_column_if_nullable(const ColumnPtr& column, |
1598 | 14.0k | const NullMap** null_map) { |
1599 | 14.0k | DORIS_CHECK(column.get() != nullptr); |
1600 | 14.0k | if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) { |
1601 | 14.0k | if (null_map != nullptr) { |
1602 | 14.0k | *null_map = &nullable_column->get_null_map_data(); |
1603 | 14.0k | } |
1604 | 14.0k | return &nullable_column->get_nested_column(); |
1605 | 14.0k | } |
1606 | 4 | return column.get(); |
1607 | 14.0k | } |
1608 | | |
1609 | | static bool _requires_parent_null_map_for_alignment(const ColumnPtr& column, |
1610 | 0 | const DataTypePtr& table_type) { |
1611 | 0 | DORIS_CHECK(column.get() != nullptr); |
1612 | 0 | DORIS_CHECK(table_type != nullptr); |
1613 | 0 | if (table_type->is_nullable()) { |
1614 | 0 | const auto& nested_type = |
1615 | 0 | assert_cast<const DataTypeNullable&>(*table_type).get_nested_type(); |
1616 | 0 | if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) { |
1617 | 0 | return _requires_parent_null_map_for_alignment( |
1618 | 0 | nullable_column->get_nested_column_ptr(), nested_type); |
1619 | 0 | } |
1620 | 0 | return _requires_parent_null_map_for_alignment(column, nested_type); |
1621 | 0 | } |
1622 | 0 | if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) { |
1623 | 0 | if (nullable_column->has_null()) { |
1624 | 0 | return true; |
1625 | 0 | } |
1626 | 0 | return _requires_parent_null_map_for_alignment(nullable_column->get_nested_column_ptr(), |
1627 | 0 | table_type); |
1628 | 0 | } |
1629 | 0 | if (const auto* array_type = typeid_cast<const DataTypeArray*>(table_type.get())) { |
1630 | 0 | const auto& array_column = assert_cast<const ColumnArray&>(*column); |
1631 | 0 | return _requires_parent_null_map_for_alignment(array_column.get_data_ptr(), |
1632 | 0 | array_type->get_nested_type()); |
1633 | 0 | } |
1634 | 0 | if (const auto* map_type = typeid_cast<const DataTypeMap*>(table_type.get())) { |
1635 | 0 | const auto& map_column = assert_cast<const ColumnMap&>(*column); |
1636 | 0 | return _requires_parent_null_map_for_alignment(map_column.get_keys_ptr(), |
1637 | 0 | map_type->get_key_type()) || |
1638 | 0 | _requires_parent_null_map_for_alignment(map_column.get_values_ptr(), |
1639 | 0 | map_type->get_value_type()); |
1640 | 0 | } |
1641 | 0 | if (const auto* struct_type = typeid_cast<const DataTypeStruct*>(table_type.get())) { |
1642 | 0 | const auto& struct_column = assert_cast<const ColumnStruct&>(*column); |
1643 | 0 | DORIS_CHECK(struct_column.tuple_size() == struct_type->get_elements().size()); |
1644 | 0 | for (size_t i = 0; i < struct_column.tuple_size(); ++i) { |
1645 | 0 | if (_requires_parent_null_map_for_alignment(struct_column.get_column_ptr(i), |
1646 | 0 | struct_type->get_element(i))) { |
1647 | 0 | return true; |
1648 | 0 | } |
1649 | 0 | } |
1650 | 0 | } |
1651 | 0 | return false; |
1652 | 0 | } |
1653 | | |
1654 | | static bool _requires_parent_null_map_for_alignment_at(const ColumnPtr& column, |
1655 | | const DataTypePtr& table_type, |
1656 | 17 | const size_t row) { |
1657 | 17 | DORIS_CHECK(column.get() != nullptr); |
1658 | 17 | DORIS_CHECK(table_type != nullptr); |
1659 | 17 | DORIS_CHECK(row < column->size()); |
1660 | 17 | if (table_type->is_nullable()) { |
1661 | 5 | const auto& nested_type = |
1662 | 5 | assert_cast<const DataTypeNullable&>(*table_type).get_nested_type(); |
1663 | 5 | if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) { |
1664 | | // A nearer nullable wrapper already protects its descendants at this entry, so an |
1665 | | // inherited collection mask cannot be needed there. |
1666 | 5 | if (nullable_column->is_null_at(row)) { |
1667 | 0 | return false; |
1668 | 0 | } |
1669 | 5 | return _requires_parent_null_map_for_alignment_at( |
1670 | 5 | nullable_column->get_nested_column_ptr(), nested_type, row); |
1671 | 5 | } |
1672 | 0 | return _requires_parent_null_map_for_alignment_at(column, nested_type, row); |
1673 | 5 | } |
1674 | 12 | if (const auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) { |
1675 | 4 | if (nullable_column->is_null_at(row)) { |
1676 | 2 | return true; |
1677 | 2 | } |
1678 | 2 | return _requires_parent_null_map_for_alignment_at( |
1679 | 2 | nullable_column->get_nested_column_ptr(), table_type, row); |
1680 | 4 | } |
1681 | 8 | if (const auto* array_type = typeid_cast<const DataTypeArray*>(table_type.get())) { |
1682 | 0 | const auto& array_column = assert_cast<const ColumnArray&>(*column); |
1683 | 0 | const auto& offsets = array_column.get_offsets(); |
1684 | 0 | const size_t begin = row == 0 ? 0 : offsets[row - 1]; |
1685 | 0 | const size_t end = offsets[row]; |
1686 | 0 | for (size_t child_row = begin; child_row < end; ++child_row) { |
1687 | 0 | if (_requires_parent_null_map_for_alignment_at(array_column.get_data_ptr(), |
1688 | 0 | array_type->get_nested_type(), |
1689 | 0 | child_row)) { |
1690 | 0 | return true; |
1691 | 0 | } |
1692 | 0 | } |
1693 | 0 | return false; |
1694 | 0 | } |
1695 | 8 | if (const auto* map_type = typeid_cast<const DataTypeMap*>(table_type.get())) { |
1696 | 0 | const auto& map_column = assert_cast<const ColumnMap&>(*column); |
1697 | 0 | const auto& offsets = map_column.get_offsets(); |
1698 | 0 | const size_t begin = row == 0 ? 0 : offsets[row - 1]; |
1699 | 0 | const size_t end = offsets[row]; |
1700 | 0 | for (size_t child_row = begin; child_row < end; ++child_row) { |
1701 | 0 | if (_requires_parent_null_map_for_alignment_at( |
1702 | 0 | map_column.get_keys_ptr(), map_type->get_key_type(), child_row) || |
1703 | 0 | _requires_parent_null_map_for_alignment_at( |
1704 | 0 | map_column.get_values_ptr(), map_type->get_value_type(), child_row)) { |
1705 | 0 | return true; |
1706 | 0 | } |
1707 | 0 | } |
1708 | 0 | return false; |
1709 | 0 | } |
1710 | 8 | if (const auto* struct_type = typeid_cast<const DataTypeStruct*>(table_type.get())) { |
1711 | 4 | const auto& struct_column = assert_cast<const ColumnStruct&>(*column); |
1712 | 4 | DORIS_CHECK(struct_column.tuple_size() == struct_type->get_elements().size()); |
1713 | 6 | for (size_t i = 0; i < struct_column.tuple_size(); ++i) { |
1714 | 4 | if (_requires_parent_null_map_for_alignment_at(struct_column.get_column_ptr(i), |
1715 | 4 | struct_type->get_element(i), row)) { |
1716 | 2 | return true; |
1717 | 2 | } |
1718 | 4 | } |
1719 | 4 | } |
1720 | 6 | return false; |
1721 | 8 | } |
1722 | | |
1723 | | static bool _requires_collection_parent_null_map(const NullMap* parent_null_map, |
1724 | | const ColumnPtr& column, |
1725 | | const DataTypePtr& table_type) { |
1726 | | // Descendant null maps can be entry-sized. Scan them only when an inherited mask can |
1727 | | // actually hide a row; absent/all-clear masks cannot authorize any physical child NULL. |
1728 | | if (parent_null_map == nullptr || |
1729 | | std::ranges::none_of(*parent_null_map, [](const auto value) { return value != 0; })) { |
1730 | | return false; |
1731 | | } |
1732 | | return _requires_parent_null_map_for_alignment(column, table_type); |
1733 | | } |
1734 | | |
1735 | | template <typename Offsets> |
1736 | | static bool _parent_null_map_hides_collection_entries(const NullMap* container_null_map, |
1737 | | const NullMap* ancestor_null_map, |
1738 | | const size_t rows, |
1739 | 9.23k | const Offsets& offsets) { |
1740 | 9.23k | if (container_null_map == nullptr && ancestor_null_map == nullptr) { |
1741 | 2 | return false; |
1742 | 2 | } |
1743 | 9.22k | DORIS_CHECK(container_null_map == nullptr || container_null_map->size() == rows); |
1744 | 9.22k | DORIS_CHECK(ancestor_null_map == nullptr || ancestor_null_map->size() == rows); |
1745 | 9.22k | DORIS_CHECK(offsets.size() == rows); |
1746 | 9.22k | size_t begin = 0; |
1747 | 21.8k | for (size_t row = 0; row < rows; ++row) { |
1748 | 12.5k | const size_t end = offsets[row]; |
1749 | 12.5k | const bool hidden = (container_null_map != nullptr && (*container_null_map)[row]) || |
1750 | 12.5k | (ancestor_null_map != nullptr && (*ancestor_null_map)[row]); |
1751 | | // A hidden collection row protects descendants only when its offset span is nonempty. |
1752 | 12.5k | if (hidden && end > begin) { |
1753 | 5 | return true; |
1754 | 5 | } |
1755 | 12.5k | begin = end; |
1756 | 12.5k | } |
1757 | 9.22k | return false; |
1758 | 9.22k | } |
1759 | | |
1760 | | template <typename Offsets> |
1761 | | static bool _requires_collection_parent_null_map(const NullMap* parent_null_map, |
1762 | | const ColumnPtr& column, |
1763 | | const DataTypePtr& table_type, |
1764 | 597 | const size_t rows, const Offsets& offsets) { |
1765 | 597 | if (parent_null_map == nullptr) { |
1766 | 0 | return false; |
1767 | 0 | } |
1768 | 597 | DORIS_CHECK(parent_null_map->size() == rows); |
1769 | 597 | DORIS_CHECK(offsets.size() == rows); |
1770 | 597 | DORIS_CHECK(offsets.empty() || offsets.back() == column->size()); |
1771 | 597 | size_t begin = 0; |
1772 | 3.44k | for (size_t row = 0; row < rows; ++row) { |
1773 | 2.85k | const size_t end = offsets[row]; |
1774 | 2.85k | if ((*parent_null_map)[row]) { |
1775 | | // Only ancestor-hidden entries can consume this projection. Restricting the probe |
1776 | | // to their spans avoids scanning visible payload covered by nearer nullable masks. |
1777 | 746 | for (size_t child_row = begin; child_row < end; ++child_row) { |
1778 | 6 | if (_requires_parent_null_map_for_alignment_at(column, table_type, child_row)) { |
1779 | 2 | return true; |
1780 | 2 | } |
1781 | 6 | } |
1782 | 742 | } |
1783 | 2.84k | begin = end; |
1784 | 2.84k | } |
1785 | 595 | return false; |
1786 | 597 | } |
1787 | | |
1788 | | template <typename Offsets> |
1789 | | static const NullMap* _project_collection_parent_null_map_for_hidden_entries( |
1790 | | const NullMap* container_null_map, const NullMap* ancestor_null_map, const size_t rows, |
1791 | 9.23k | const Offsets& offsets, const size_t child_rows, NullMap* const projected_null_map) { |
1792 | 9.23k | if (!_parent_null_map_hides_collection_entries(container_null_map, ancestor_null_map, rows, |
1793 | 9.23k | offsets)) { |
1794 | | // Nullable collection wrappers expose a null-map even when every row is present; avoid |
1795 | | // allocating entry-coordinate scratch unless a hidden row owns physical entries. |
1796 | 9.22k | return nullptr; |
1797 | 9.22k | } |
1798 | 7 | projected_null_map->resize(child_rows); |
1799 | 7 | std::fill(projected_null_map->begin(), projected_null_map->end(), 0); |
1800 | 7 | size_t begin = 0; |
1801 | 17 | for (size_t row = 0; row < rows; ++row) { |
1802 | 10 | const size_t end = offsets[row]; |
1803 | 10 | const bool hidden = (container_null_map != nullptr && (*container_null_map)[row]) || |
1804 | 10 | (ancestor_null_map != nullptr && (*ancestor_null_map)[row]); |
1805 | 10 | if (hidden) { |
1806 | | // Collection masks use row coordinates; descendants need the same invariant |
1807 | | // projected through offsets so hidden physical payload cannot fail validation. |
1808 | 5 | std::fill(projected_null_map->begin() + begin, projected_null_map->begin() + end, |
1809 | 5 | 1); |
1810 | 5 | } |
1811 | 10 | begin = end; |
1812 | 10 | } |
1813 | 7 | DORIS_CHECK(begin == child_rows); |
1814 | 7 | return projected_null_map; |
1815 | 9.23k | } |
1816 | | |
1817 | | Status _materialize_struct_mapping_column(const ColumnMapping& mapping, |
1818 | | const ColumnPtr& file_column, const size_t rows, |
1819 | | ColumnPtr* column, |
1820 | 4.79k | const NullMap* nullable_parent_null_map = nullptr) { |
1821 | 4.79k | DORIS_CHECK(mapping.table_type != nullptr); |
1822 | 4.79k | const auto* table_type = |
1823 | 4.79k | assert_cast<const DataTypeStruct*>(remove_nullable(mapping.table_type).get()); |
1824 | 4.79k | const auto full_file_column = file_column->convert_to_full_column_if_const(); |
1825 | 4.79k | const NullMap* parent_null_map = nullptr; |
1826 | 4.79k | const auto* nested_file_column = |
1827 | 4.79k | _nested_column_if_nullable(full_file_column, &parent_null_map); |
1828 | 4.79k | const auto* file_struct = assert_cast<const ColumnStruct*>(nested_file_column); |
1829 | 4.79k | DORIS_CHECK(table_type->get_elements().size() == mapping.child_mappings.size()); |
1830 | | |
1831 | 4.79k | NullMap combined_parent_null_map; |
1832 | 4.79k | const NullMap* descendant_parent_null_map = nullable_parent_null_map; |
1833 | 4.79k | if (parent_null_map != nullptr) { |
1834 | 4.79k | DORIS_CHECK(parent_null_map->size() == rows); |
1835 | 4.79k | if (nullable_parent_null_map != nullptr) { |
1836 | 51 | DORIS_CHECK(nullable_parent_null_map->size() == rows); |
1837 | 51 | } |
1838 | 4.79k | if (!mapping.table_type->is_nullable()) { |
1839 | 7 | for (size_t i = 0; i < rows; ++i) { |
1840 | | // A required nested container may drop its own NULL only when an ancestor |
1841 | | // already hides that row; otherwise physical defaults become visible values. |
1842 | 5 | if ((*parent_null_map)[i] && |
1843 | 5 | (nullable_parent_null_map == nullptr || !(*nullable_parent_null_map)[i])) { |
1844 | 1 | return Status::InternalError( |
1845 | 1 | "Source struct contains NULL for non-nullable table column"); |
1846 | 1 | } |
1847 | 5 | } |
1848 | 3 | } |
1849 | 4.79k | combined_parent_null_map.resize(rows); |
1850 | 104k | for (size_t i = 0; i < rows; ++i) { |
1851 | 99.7k | combined_parent_null_map[i] = |
1852 | 99.7k | (*parent_null_map)[i] || |
1853 | 99.7k | (nullable_parent_null_map != nullptr && (*nullable_parent_null_map)[i]); |
1854 | 99.7k | } |
1855 | 4.79k | descendant_parent_null_map = &combined_parent_null_map; |
1856 | 4.79k | } |
1857 | | |
1858 | 4.79k | Columns child_columns; |
1859 | 4.79k | child_columns.reserve(mapping.child_mappings.size()); |
1860 | 4.79k | const auto file_ordered_children = |
1861 | 4.79k | _present_child_mappings_in_file_order(mapping.child_mappings); |
1862 | 4.79k | const auto table_ordered_children = |
1863 | 4.79k | _child_mappings_in_table_type_order(mapping, *table_type); |
1864 | 12.7k | for (const auto* child_mapping : table_ordered_children) { |
1865 | 12.7k | DORIS_CHECK(child_mapping != nullptr); |
1866 | 12.7k | if (!child_mapping->file_local_id.has_value()) { |
1867 | 5.02k | ColumnPtr child_column; |
1868 | 5.02k | RETURN_IF_ERROR(_materialize_default_or_missing_column( |
1869 | 5.02k | *child_mapping, nullptr, rows, &child_column, descendant_parent_null_map)); |
1870 | 5.02k | child_column = child_column->convert_to_full_column_if_const(); |
1871 | 5.02k | child_columns.push_back(std::move(child_column)); |
1872 | 5.02k | continue; |
1873 | 5.02k | } |
1874 | 7.73k | const auto file_child_idx = |
1875 | 7.73k | _file_child_ordinal_for_mapping(mapping, *child_mapping, file_ordered_children); |
1876 | 7.73k | DORIS_CHECK(file_child_idx < file_struct->get_columns().size()); |
1877 | 7.73k | ColumnPtr child_column = file_struct->get_column_ptr(file_child_idx); |
1878 | 7.73k | RETURN_IF_ERROR(_materialize_present_child_mapping_column( |
1879 | 7.73k | *child_mapping, child_column, rows, &child_column, descendant_parent_null_map)); |
1880 | 7.73k | child_columns.push_back(std::move(child_column)); |
1881 | 7.73k | } |
1882 | 4.79k | MutableColumns mutable_child_columns; |
1883 | 4.79k | mutable_child_columns.reserve(child_columns.size()); |
1884 | 12.7k | for (auto& child_column : child_columns) { |
1885 | 12.7k | mutable_child_columns.push_back(IColumn::mutate(std::move(child_column))); |
1886 | 12.7k | } |
1887 | 4.79k | auto result = ColumnStruct::create(std::move(mutable_child_columns)); |
1888 | 4.79k | if (mapping.table_type->is_nullable()) { |
1889 | 4.78k | auto null_map = ColumnUInt8::create(); |
1890 | 4.78k | auto& null_map_data = null_map->get_data(); |
1891 | 4.78k | null_map_data.resize(rows); |
1892 | 4.78k | if (parent_null_map != nullptr) { |
1893 | 4.78k | DORIS_CHECK(parent_null_map->size() == rows); |
1894 | 4.78k | null_map_data.assign(parent_null_map->begin(), parent_null_map->end()); |
1895 | 4.78k | } else { |
1896 | 0 | std::fill(null_map_data.begin(), null_map_data.end(), 0); |
1897 | 0 | } |
1898 | 4.78k | *column = ColumnNullable::create(std::move(result), std::move(null_map)); |
1899 | 4.78k | } else { |
1900 | 4 | *column = std::move(result); |
1901 | 4 | } |
1902 | 4.79k | return Status::OK(); |
1903 | 4.79k | } |
1904 | | |
1905 | | Status _materialize_array_mapping_column(const ColumnMapping& mapping, |
1906 | | const ColumnPtr& file_column, const size_t rows, |
1907 | | ColumnPtr* column, |
1908 | 4.83k | const NullMap* nullable_parent_null_map = nullptr) { |
1909 | 4.83k | DORIS_CHECK(mapping.child_mappings.size() == 1); |
1910 | 4.83k | const auto full_file_column = file_column->convert_to_full_column_if_const(); |
1911 | 4.83k | const NullMap* parent_null_map = nullptr; |
1912 | 4.83k | const auto* nested_file_column = |
1913 | 4.83k | _nested_column_if_nullable(full_file_column, &parent_null_map); |
1914 | 4.83k | if (parent_null_map != nullptr && !mapping.table_type->is_nullable()) { |
1915 | 2 | DORIS_CHECK(parent_null_map->size() == rows); |
1916 | 2 | if (nullable_parent_null_map != nullptr) { |
1917 | 1 | DORIS_CHECK(nullable_parent_null_map->size() == rows); |
1918 | 1 | } |
1919 | 4 | for (size_t i = 0; i < rows; ++i) { |
1920 | | // ARRAY row masks cannot be forwarded to elements because they use different |
1921 | | // coordinates, so validate the container before dropping its nullable wrapper. |
1922 | 3 | if ((*parent_null_map)[i] && |
1923 | 3 | (nullable_parent_null_map == nullptr || !(*nullable_parent_null_map)[i])) { |
1924 | 1 | return Status::InternalError( |
1925 | 1 | "Source array contains NULL for non-nullable table column"); |
1926 | 1 | } |
1927 | 3 | } |
1928 | 2 | } |
1929 | 4.83k | const auto* file_array = assert_cast<const ColumnArray*>(nested_file_column); |
1930 | 4.83k | ColumnPtr nested_column = file_array->get_data_ptr(); |
1931 | 4.83k | auto element_mapping = mapping.child_mappings[0]; |
1932 | | // Keep the descriptor type for schema matching. ARRAY's nullable element wrapper is a |
1933 | | // storage invariant, so add it only at the materialization boundary. |
1934 | 4.83k | element_mapping.table_type = make_nullable(element_mapping.table_type); |
1935 | 4.83k | NullMap descendant_parent_null_map; |
1936 | 4.83k | const NullMap* descendant_parent_null_map_ptr = |
1937 | 4.83k | _project_collection_parent_null_map_for_hidden_entries( |
1938 | 4.83k | parent_null_map, nullable_parent_null_map, rows, file_array->get_offsets(), |
1939 | 4.83k | nested_column->size(), &descendant_parent_null_map); |
1940 | 4.83k | RETURN_IF_ERROR(_materialize_present_child_mapping_column( |
1941 | 4.83k | element_mapping, nested_column, nested_column->size(), &nested_column, |
1942 | 4.83k | descendant_parent_null_map_ptr)); |
1943 | 4.83k | auto offsets_column = file_array->get_offsets_ptr()->convert_to_full_column_if_const(); |
1944 | 4.83k | auto result = ColumnArray::create(IColumn::mutate(std::move(nested_column)), |
1945 | 4.83k | IColumn::mutate(std::move(offsets_column))); |
1946 | 4.83k | if (mapping.table_type->is_nullable()) { |
1947 | 4.83k | auto null_map = ColumnUInt8::create(); |
1948 | 4.83k | auto& null_map_data = null_map->get_data(); |
1949 | 4.83k | null_map_data.resize(rows); |
1950 | 4.83k | if (parent_null_map != nullptr) { |
1951 | 4.83k | DORIS_CHECK(parent_null_map->size() == rows); |
1952 | 4.83k | null_map_data.assign(parent_null_map->begin(), parent_null_map->end()); |
1953 | 4.83k | } else { |
1954 | 0 | std::fill(null_map_data.begin(), null_map_data.end(), 0); |
1955 | 0 | } |
1956 | 4.83k | *column = ColumnNullable::create(std::move(result), std::move(null_map)); |
1957 | 4.83k | } else { |
1958 | 1 | *column = std::move(result); |
1959 | 1 | } |
1960 | 4.83k | return Status::OK(); |
1961 | 4.83k | } |
1962 | | |
1963 | | Status _materialize_map_mapping_column(const ColumnMapping& mapping, |
1964 | | const ColumnPtr& file_column, const size_t rows, |
1965 | | ColumnPtr* column, |
1966 | 4.39k | const NullMap* nullable_parent_null_map = nullptr) { |
1967 | 4.39k | const auto full_file_column = file_column->convert_to_full_column_if_const(); |
1968 | 4.39k | const NullMap* parent_null_map = nullptr; |
1969 | 4.39k | const auto* nested_file_column = |
1970 | 4.39k | _nested_column_if_nullable(full_file_column, &parent_null_map); |
1971 | 4.39k | if (parent_null_map != nullptr && !mapping.table_type->is_nullable()) { |
1972 | 0 | DORIS_CHECK(parent_null_map->size() == rows); |
1973 | 0 | if (nullable_parent_null_map != nullptr) { |
1974 | 0 | DORIS_CHECK(nullable_parent_null_map->size() == rows); |
1975 | 0 | } |
1976 | 0 | for (size_t i = 0; i < rows; ++i) { |
1977 | | // MAP row masks cannot be forwarded to entries because they use different |
1978 | | // coordinates, so validate the container before dropping its nullable wrapper. |
1979 | 0 | if ((*parent_null_map)[i] && |
1980 | 0 | (nullable_parent_null_map == nullptr || !(*nullable_parent_null_map)[i])) { |
1981 | 0 | return Status::InternalError( |
1982 | 0 | "Source map contains NULL for non-nullable table column"); |
1983 | 0 | } |
1984 | 0 | } |
1985 | 0 | } |
1986 | 4.39k | const auto* file_map = assert_cast<const ColumnMap*>(nested_file_column); |
1987 | 4.39k | ColumnPtr key_column = file_map->get_keys_ptr(); |
1988 | 4.39k | ColumnPtr value_column = file_map->get_values_ptr(); |
1989 | 4.39k | DORIS_CHECK(key_column->size() == value_column->size()); |
1990 | 4.39k | NullMap descendant_parent_null_map; |
1991 | 4.39k | const NullMap* descendant_parent_null_map_ptr = |
1992 | 4.39k | _project_collection_parent_null_map_for_hidden_entries( |
1993 | 4.39k | parent_null_map, nullable_parent_null_map, rows, file_map->get_offsets(), |
1994 | 4.39k | key_column->size(), &descendant_parent_null_map); |
1995 | | |
1996 | 4.39k | const ColumnMapping* key_mapping = nullptr; |
1997 | 4.39k | const ColumnMapping* value_mapping = nullptr; |
1998 | 8.78k | for (const auto& child_mapping : mapping.child_mappings) { |
1999 | 8.78k | if (!child_mapping.file_local_id.has_value()) { |
2000 | 0 | continue; |
2001 | 0 | } |
2002 | 8.78k | if (*child_mapping.file_local_id == 0) { |
2003 | 4.39k | key_mapping = &child_mapping; |
2004 | 4.39k | } else if (*child_mapping.file_local_id == 1) { |
2005 | 4.39k | value_mapping = &child_mapping; |
2006 | 4.39k | } |
2007 | 8.78k | } |
2008 | | |
2009 | 4.39k | if (key_mapping != nullptr) { |
2010 | 4.39k | RETURN_IF_ERROR(_materialize_present_child_mapping_column( |
2011 | 4.39k | *key_mapping, key_column, key_column->size(), &key_column, |
2012 | 4.39k | descendant_parent_null_map_ptr)); |
2013 | 4.39k | } else { |
2014 | 5 | const auto* table_map = |
2015 | 5 | assert_cast<const DataTypeMap*>(remove_nullable(mapping.table_type).get()); |
2016 | | // Value-only projection retains the physical key stream to preserve entry offsets; |
2017 | | // align it under the entry mask so NULL placeholders from hidden Map rows stay hidden. |
2018 | 5 | RETURN_IF_ERROR(_align_column_nullability(&key_column, table_map->get_key_type(), |
2019 | 5 | descendant_parent_null_map_ptr)); |
2020 | 5 | } |
2021 | 4.39k | if (value_mapping != nullptr) { |
2022 | 4.39k | RETURN_IF_ERROR(_materialize_present_child_mapping_column( |
2023 | 4.39k | *value_mapping, value_column, value_column->size(), &value_column, |
2024 | 4.39k | descendant_parent_null_map_ptr)); |
2025 | 4.39k | } else { |
2026 | 2 | const auto* table_map = |
2027 | 2 | assert_cast<const DataTypeMap*>(remove_nullable(mapping.table_type).get()); |
2028 | | // A retained structural value stream follows the same hidden-entry invariant as keys. |
2029 | 2 | RETURN_IF_ERROR(_align_column_nullability(&value_column, table_map->get_value_type(), |
2030 | 2 | descendant_parent_null_map_ptr)); |
2031 | 2 | } |
2032 | 4.39k | auto offsets_column = file_map->get_offsets_ptr()->convert_to_full_column_if_const(); |
2033 | 4.39k | auto result = ColumnMap::create(IColumn::mutate(std::move(key_column)), |
2034 | 4.39k | IColumn::mutate(std::move(value_column)), |
2035 | 4.39k | IColumn::mutate(std::move(offsets_column))); |
2036 | 4.39k | if (mapping.table_type->is_nullable()) { |
2037 | 4.39k | auto null_map = ColumnUInt8::create(); |
2038 | 4.39k | auto& null_map_data = null_map->get_data(); |
2039 | 4.39k | null_map_data.resize(rows); |
2040 | 4.39k | if (parent_null_map != nullptr) { |
2041 | 4.39k | DORIS_CHECK(parent_null_map->size() == rows); |
2042 | 4.39k | null_map_data.assign(parent_null_map->begin(), parent_null_map->end()); |
2043 | 4.39k | } else { |
2044 | 0 | std::fill(null_map_data.begin(), null_map_data.end(), 0); |
2045 | 0 | } |
2046 | 4.39k | *column = ColumnNullable::create(std::move(result), std::move(null_map)); |
2047 | 4.39k | } else { |
2048 | 2 | *column = std::move(result); |
2049 | 2 | } |
2050 | 4.39k | return Status::OK(); |
2051 | 4.39k | } |
2052 | | |
2053 | 855k | Status _open_mapping_expr_tree(const ColumnMapping& mapping, const RowDescriptor& row_desc) { |
2054 | 855k | if (mapping.projection != nullptr) { |
2055 | 502k | RETURN_IF_ERROR(mapping.projection->prepare(_runtime_state, row_desc)); |
2056 | 502k | RETURN_IF_ERROR(mapping.projection->open(_runtime_state)); |
2057 | 502k | } |
2058 | 855k | if (mapping.default_expr != nullptr) { |
2059 | 22.3k | RETURN_IF_ERROR(mapping.default_expr->prepare(_runtime_state, row_desc)); |
2060 | 22.3k | RETURN_IF_ERROR(mapping.default_expr->open(_runtime_state)); |
2061 | 22.3k | } |
2062 | 855k | for (const auto& child_mapping : mapping.child_mappings) { |
2063 | 329k | RETURN_IF_ERROR(_open_mapping_expr_tree(child_mapping, row_desc)); |
2064 | 329k | } |
2065 | 855k | return Status::OK(); |
2066 | 855k | } |
2067 | | |
2068 | 120k | Status _open_mapping_exprs() { |
2069 | 120k | RowDescriptor row_desc; |
2070 | 527k | for (const auto& mapping : _data_reader.column_mapper->mappings()) { |
2071 | 527k | RETURN_IF_ERROR(_open_mapping_expr_tree(mapping, row_desc)); |
2072 | 527k | } |
2073 | 120k | return Status::OK(); |
2074 | 120k | } |
2075 | | |
2076 | | Status _build_file_aggregate_request(TPushAggOp::type agg_type, |
2077 | 1.51k | FileAggregateRequest* request) const { |
2078 | 1.51k | DORIS_CHECK(request != nullptr); |
2079 | 1.51k | DORIS_CHECK(_supports_aggregate_pushdown(agg_type)); |
2080 | 1.51k | request->agg_type = agg_type; |
2081 | 1.51k | request->columns.clear(); |
2082 | 1.51k | if (agg_type == TPushAggOp::type::COUNT) { |
2083 | 1.48k | DORIS_CHECK(_push_down_count_columns.has_value()); |
2084 | | // An empty explicit list is the semantic signal for COUNT(*). Do not inspect the |
2085 | | // mapping count: `SELECT COUNT(*) FROM t` may still project one nullable column because |
2086 | | // the planner keeps a placeholder slot. In a 10,000-row file where that arbitrary slot |
2087 | | // has 9,015 non-null values, passing the slot would ask Parquet/ORC metadata for |
2088 | | // COUNT(slot)=9,015 instead of the required row count 10,000. |
2089 | 1.48k | if (!_push_down_count_columns->empty()) { |
2090 | 252 | const auto& mapping = _push_down_count_mapping(); |
2091 | 252 | DORIS_CHECK(mapping.file_local_id.has_value()); |
2092 | 252 | FileAggregateRequest::Column column; |
2093 | 252 | column.projection = |
2094 | 252 | LocalColumnIndex::top_level(LocalColumnId(*mapping.file_local_id)); |
2095 | 252 | request->columns.push_back(std::move(column)); |
2096 | 252 | } |
2097 | 1.48k | return Status::OK(); |
2098 | 1.48k | } |
2099 | 28 | request->columns.reserve(_data_reader.column_mapper->mappings().size()); |
2100 | 47 | for (const auto& mapping : _data_reader.column_mapper->mappings()) { |
2101 | 47 | DORIS_CHECK(mapping.file_local_id.has_value()); |
2102 | 47 | FileAggregateRequest::Column column; |
2103 | 47 | column.projection = LocalColumnIndex::top_level(LocalColumnId(*mapping.file_local_id)); |
2104 | 47 | if (!mapping.child_mappings.empty()) { |
2105 | 1 | RETURN_IF_ERROR(build_aggregate_projection(mapping, &column.projection)); |
2106 | 1 | } |
2107 | 47 | request->columns.push_back(std::move(column)); |
2108 | 47 | } |
2109 | 28 | return Status::OK(); |
2110 | 28 | } |
2111 | | |
2112 | 1.03k | const ColumnMapping& _push_down_count_mapping() const { |
2113 | 1.03k | DORIS_CHECK(_push_down_count_columns.has_value()); |
2114 | 1.03k | DORIS_CHECK(_push_down_count_columns->size() == 1); |
2115 | 1.03k | const auto mapping_it = |
2116 | 1.03k | std::ranges::find(_data_reader.column_mapper->mappings(), |
2117 | 1.03k | _push_down_count_columns->front(), &ColumnMapping::global_index); |
2118 | | // FileScannerV2 translates FE SlotIds through the same projected-column list used to build |
2119 | | // the mapper, so a missing mapping is an FE/BE contract violation rather than a fallback. |
2120 | 1.03k | DORIS_CHECK(mapping_it != _data_reader.column_mapper->mappings().end()); |
2121 | 1.03k | return *mapping_it; |
2122 | 1.03k | } |
2123 | | |
2124 | | Status _materialize_aggregate_pushdown_rows(TPushAggOp::type agg_type, |
2125 | | const FileAggregateResult& file_result, |
2126 | 21 | Block* block) { |
2127 | 21 | DORIS_CHECK(agg_type == TPushAggOp::type::MINMAX); |
2128 | | // MIN/MAX pushdown emits two rows, min first and max second, for each projected column. |
2129 | | // The upper MIN/MAX aggregate consumes those two rows to produce the final aggregate value. |
2130 | 21 | DORIS_CHECK(file_result.columns.size() == _data_reader.column_mapper->mappings().size()); |
2131 | 21 | DORIS_CHECK(block->columns() == _data_reader.column_mapper->mappings().size()); |
2132 | 21 | Block file_block; |
2133 | 21 | file_block.reserve(_data_reader.file_block_layout.size()); |
2134 | 26 | for (const auto& column : _data_reader.file_block_layout) { |
2135 | 26 | file_block.insert({column.type->create_column(), column.type, column.name}); |
2136 | 26 | } |
2137 | 47 | for (size_t column_idx = 0; column_idx < file_result.columns.size(); ++column_idx) { |
2138 | 26 | const auto& result_column = file_result.columns[column_idx]; |
2139 | 26 | if (!result_column.has_min || !result_column.has_max) { |
2140 | 0 | return Status::NotSupported("Missing min/max aggregate result for column {}", |
2141 | 0 | _projected_columns[column_idx].name); |
2142 | 0 | } |
2143 | 26 | bool found_file_column = false; |
2144 | 33 | for (size_t block_position = 0; block_position < _data_reader.file_block_layout.size(); |
2145 | 33 | ++block_position) { |
2146 | 33 | if (_data_reader.file_block_layout[block_position].file_column_id == |
2147 | 33 | file_result.columns[column_idx].projection.column_id()) { |
2148 | 26 | found_file_column = true; |
2149 | 26 | auto column = file_block.get_by_position(block_position) |
2150 | 26 | .type->create_column() |
2151 | 26 | ->assert_mutable(); |
2152 | 26 | RETURN_IF_ERROR(_insert_aggregate_projection_value( |
2153 | 26 | file_result.columns[column_idx].projection, result_column.min_value, |
2154 | 26 | column.get())); |
2155 | 26 | RETURN_IF_ERROR(_insert_aggregate_projection_value( |
2156 | 26 | file_result.columns[column_idx].projection, result_column.max_value, |
2157 | 26 | column.get())); |
2158 | 26 | file_block.replace_by_position(block_position, std::move(column)); |
2159 | 26 | break; |
2160 | 26 | } |
2161 | 33 | } |
2162 | 26 | DORIS_CHECK(found_file_column); |
2163 | 26 | } |
2164 | 47 | for (size_t column_idx = 0; column_idx < _data_reader.column_mapper->mappings().size(); |
2165 | 26 | ++column_idx) { |
2166 | 26 | ColumnPtr table_column; |
2167 | 26 | RETURN_IF_ERROR(_materialize_mapping_column( |
2168 | 26 | _data_reader.column_mapper->mappings()[column_idx], &file_block, 2, |
2169 | 26 | &table_column, |
2170 | 26 | column_idx + 1 == _data_reader.column_mapper->mappings().size())); |
2171 | 26 | block->replace_by_position(column_idx, std::move(table_column)); |
2172 | 26 | } |
2173 | 21 | return Status::OK(); |
2174 | 21 | } |
2175 | | |
2176 | | struct FileBlockColumn { |
2177 | | LocalColumnId file_column_id = LocalColumnId::invalid(); |
2178 | | std::string name; |
2179 | | DataTypePtr type; |
2180 | | }; |
2181 | | |
2182 | | struct DataReader { |
2183 | | std::unique_ptr<FileReader> reader; |
2184 | | std::unique_ptr<TableColumnMapper> column_mapper; |
2185 | | // Schema of the data file, also including virtual column (row position). |
2186 | | std::vector<ColumnDefinition> file_schema; |
2187 | | // Layout of the block returned by file reader, determined by column mapping and file |
2188 | | // schema. It is used for file reader to materialize columns into correct type and position. |
2189 | | std::vector<FileBlockColumn> file_block_layout; |
2190 | | Block block_template; |
2191 | | }; |
2192 | | DataReader _data_reader; |
2193 | | // Latest immutable request queued to the physical reader. The file-block layout remains fixed |
2194 | | // for the split even while predicates are refreshed at a reader-defined granule boundary. |
2195 | | std::shared_ptr<FileScanRequest> _file_scan_request; |
2196 | | std::vector<ColumnDefinition> _projected_columns; |
2197 | | std::unique_ptr<ScanTask> _current_task; |
2198 | | std::optional<io::FileDescription> _current_file_description; |
2199 | | // Range-level compression has higher priority than scan-param compression. TVF/load can keep |
2200 | | // the logical format as CSV/TEXT while carrying the concrete compression such as GZ or LZO on |
2201 | | // each TFileRangeDesc, matching the old FileScanner reader contract. |
2202 | | TFileCompressType::type _current_range_compress_type = TFileCompressType::UNKNOWN; |
2203 | | std::optional<TUniqueId> _current_range_load_id; |
2204 | | TFileRangeDesc _current_file_range_desc; |
2205 | | std::shared_ptr<io::FileSystemProperties> _system_properties; |
2206 | | // partition key -> value |
2207 | | std::map<std::string, Field> _partition_values; |
2208 | | // Predicates built from scan conjuncts before file-level localization. |
2209 | | std::vector<TableFilter> _table_filters; |
2210 | | // Number of localized filters before the first unsafe conjunct in the original row-level |
2211 | | // order. This differs from scanning `_table_filters` for safety because slotless predicates are |
2212 | | // intentionally absent from that vector but must still act as ordering barriers. |
2213 | | size_t _constant_pruning_safe_filter_count = 0; |
2214 | | VExprContextSPtrs _conjuncts; |
2215 | | ReadProfile _profile; |
2216 | | // Parsed from row-position based delete files, including position delete and deletion vector. |
2217 | | DeleteRows* _delete_rows = nullptr; |
2218 | | DeletionVector* _deletion_vector = nullptr; |
2219 | | TFileScanRangeParams* _scan_params; |
2220 | | std::shared_ptr<io::IOContext> _io_ctx; |
2221 | | RuntimeState* _runtime_state; |
2222 | | RuntimeProfile* _scanner_profile; |
2223 | | const std::vector<SlotDescriptor*>* _file_slot_descs = nullptr; |
2224 | | FileFormat _format; |
2225 | | TPushAggOp::type _push_down_agg_type = TPushAggOp::type::NONE; |
2226 | | std::optional<std::vector<GlobalIndex>> _push_down_count_columns; |
2227 | | size_t _batch_size = 0; |
2228 | | uint64_t _initial_condition_cache_digest = 0; |
2229 | | uint64_t _condition_cache_digest = 0; |
2230 | | // True only when prepare_split() received a digest for the exact conjunct snapshot used by |
2231 | | // this split. Standalone callers that only supplied TableReadOptions::condition_cache_digest |
2232 | | // keep the conservative runtime-filter guard. |
2233 | | bool _condition_cache_digest_covers_current_split = false; |
2234 | | segment_v2::ConditionCache::ExternalCacheKey _condition_cache_key; |
2235 | | std::shared_ptr<std::vector<bool>> _condition_cache; |
2236 | | std::shared_ptr<ConditionCacheContext> _condition_cache_ctx; |
2237 | | int64_t _condition_cache_hit_count = 0; |
2238 | | bool _current_reader_reached_eof = false; |
2239 | | int64_t _remaining_table_level_count = -1; |
2240 | | int64_t _remaining_file_level_count = -1; |
2241 | | // True only after the active split selects a table-level row-count shortcut or successfully |
2242 | | // materializes COUNT rows from file metadata. FileScannerV2 uses this result, rather than the |
2243 | | // raw aggregate opcode, to keep adaptive batching enabled for normal row-scan fallbacks. |
2244 | | bool _current_split_uses_metadata_count = false; |
2245 | | // Snapshot supplied by FileScannerV2 for the active split. It gates every shortcut that emits |
2246 | | // irreversible aggregate rows, not only the table-level row-count shortcut in prepare_split(). |
2247 | | bool _all_runtime_filters_applied_for_split = true; |
2248 | | std::optional<GlobalRowIdContext> _global_rowid_context; |
2249 | | bool _aggregate_pushdown_tried = false; |
2250 | | bool _current_split_pruned = false; |
2251 | | TableColumnMapperOptions _mapper_options; |
2252 | | |
2253 | | private: |
2254 | | static const ColumnDefinition* _find_column_definition( |
2255 | 562k | const std::vector<ColumnDefinition>& schema, LocalColumnId column_id) { |
2256 | 10.0M | for (const auto& field : schema) { |
2257 | 10.0M | if (field.file_local_id() == column_id.value()) { |
2258 | 533k | return &field; |
2259 | 533k | } |
2260 | 10.0M | } |
2261 | 28.9k | return nullptr; |
2262 | 562k | } |
2263 | | |
2264 | 142 | static bool _can_push_down_minmax_for_mapping(const ColumnMapping& mapping) { |
2265 | 142 | if (mapping.child_mappings.empty()) { |
2266 | | // Direct mappings use a slot-ref projection to materialize the file column. The |
2267 | | // projection does not transform ordering; casts and other conversions are already |
2268 | | // represented by a non-trivial mapping and must fall back to row scanning. |
2269 | 139 | return mapping.is_trivial; |
2270 | 139 | } |
2271 | 3 | const auto primitive_type = remove_nullable(mapping.file_type)->get_primitive_type(); |
2272 | 3 | if (primitive_type != TYPE_STRUCT) { |
2273 | 1 | return false; |
2274 | 1 | } |
2275 | 2 | size_t mapped_children = 0; |
2276 | 2 | const ColumnMapping* mapped_child = nullptr; |
2277 | 2 | for (const auto& child_mapping : mapping.child_mappings) { |
2278 | 2 | if (!child_mapping.file_local_id.has_value()) { |
2279 | 0 | continue; |
2280 | 0 | } |
2281 | 2 | ++mapped_children; |
2282 | 2 | mapped_child = &child_mapping; |
2283 | 2 | } |
2284 | 2 | return mapped_children == 1 && mapped_child != nullptr && |
2285 | 2 | _can_push_down_minmax_for_mapping(*mapped_child); |
2286 | 3 | } |
2287 | | |
2288 | | static Status build_aggregate_projection(const ColumnMapping& mapping, |
2289 | 2 | LocalColumnIndex* projection) { |
2290 | 2 | DORIS_CHECK(projection != nullptr); |
2291 | 2 | DORIS_CHECK(mapping.file_local_id.has_value()); |
2292 | 2 | *projection = LocalColumnIndex::local(*mapping.file_local_id); |
2293 | 2 | projection->children.clear(); |
2294 | 2 | projection->project_all_children = true; |
2295 | 2 | if (mapping.child_mappings.empty()) { |
2296 | 1 | return Status::OK(); |
2297 | 1 | } |
2298 | 1 | projection->project_all_children = false; |
2299 | 1 | for (const auto& child_mapping : mapping.child_mappings) { |
2300 | 1 | if (!child_mapping.file_local_id.has_value()) { |
2301 | 0 | continue; |
2302 | 0 | } |
2303 | 1 | LocalColumnIndex child_projection; |
2304 | 1 | RETURN_IF_ERROR(build_aggregate_projection(child_mapping, &child_projection)); |
2305 | 1 | projection->children.push_back(std::move(child_projection)); |
2306 | 1 | } |
2307 | 1 | DORIS_CHECK(projection->children.size() == 1); |
2308 | 1 | return Status::OK(); |
2309 | 1 | } |
2310 | | |
2311 | | static Status _insert_aggregate_projection_value(const LocalColumnIndex& projection, |
2312 | 108 | const Field& value, IColumn* column) { |
2313 | 108 | DORIS_CHECK(column != nullptr); |
2314 | 108 | if (auto* nullable_column = check_and_get_column<ColumnNullable>(*column)) { |
2315 | 54 | RETURN_IF_ERROR(_insert_aggregate_projection_value( |
2316 | 54 | projection, value, &nullable_column->get_nested_column())); |
2317 | 54 | nullable_column->get_null_map_data().push_back(0); |
2318 | 54 | return Status::OK(); |
2319 | 54 | } |
2320 | 54 | if (projection.project_all_children || projection.children.empty()) { |
2321 | 52 | column->insert(value); |
2322 | 52 | return Status::OK(); |
2323 | 52 | } |
2324 | 2 | auto* struct_column = assert_cast<ColumnStruct*>(column); |
2325 | 2 | DORIS_CHECK(projection.children.size() == 1); |
2326 | 2 | const auto& child_projection = projection.children[0]; |
2327 | 2 | DORIS_CHECK(struct_column->get_columns().size() == 1); |
2328 | 2 | RETURN_IF_ERROR(_insert_aggregate_projection_value(child_projection, value, |
2329 | 2 | &struct_column->get_column(0))); |
2330 | 2 | return Status::OK(); |
2331 | 2 | } |
2332 | | |
2333 | | // Parse a DV into its compressed bitmap. Position delete files continue to use _delete_rows. |
2334 | | Status _parse_delete_predicates(const SplitReadOptions& options); |
2335 | | }; |
2336 | | |
2337 | | } // namespace doris::format |