github-actions[bot] commented on code in PR #65674:
URL: https://github.com/apache/doris/pull/65674#discussion_r3610731624
##########
be/src/core/data_type_serde/data_type_time_serde.cpp:
##########
@@ -51,6 +55,64 @@ TimeValue::TimeType read_time_decoded_value(const
DecodedColumnView& view, int64
abs_micros % TimeValue::ONE_SECOND_MICROSECONDS, negative);
}
+class TimeV2ParquetConsumer final : public ParquetFixedValueConsumer {
+public:
+ TimeV2ParquetConsumer(IColumn& column, const ParquetDecodeContext& context,
+ ParquetMaterializationState* state = nullptr)
+ : _data(assert_cast<ColumnTimeV2&>(column).get_data()),
+ _context(context),
+ _state(state) {}
+
+ Status consume(const uint8_t* values, size_t num_values, size_t
value_width) override {
+ const size_t old_size = _data.size();
+ _data.resize(old_size + num_values);
+ for (size_t row = 0; row < num_values; ++row) {
+ int64_t micros;
+ if (_context.physical_type == ParquetPhysicalType::INT32) {
+ DORIS_CHECK_EQ(value_width, sizeof(int32_t));
+ micros = static_cast<int64_t>(
+ unaligned_load<int32_t>(values + row *
sizeof(int32_t))) *
+ 1000;
+ } else {
+ DORIS_CHECK(_context.physical_type ==
ParquetPhysicalType::INT64);
+ DORIS_CHECK_EQ(value_width, sizeof(int64_t));
+ micros = unaligned_load<int64_t>(values + row *
sizeof(int64_t));
+ if (_context.time_unit == ParquetTimeUnit::MILLIS) {
+ if (micros > std::numeric_limits<int64_t>::max() / 1000 ||
+ micros < std::numeric_limits<int64_t>::min() / 1000) {
+ if (_state != nullptr &&
_state->mark_conversion_failure(old_size + row)) {
+ _data[old_size + row] = TimeValue::TimeType();
+ continue;
+ }
+ _data.resize(old_size);
+ return Status::DataQualityError(
+ "Parquet TIME value overflows microseconds");
+ }
+ micros *= 1000;
+ } else if (_context.time_unit == ParquetTimeUnit::NANOS) {
+ micros /= 1000;
+ }
+ }
+ // Doris TIMEV2 stores signed microseconds in a double. Splitting
into calendar fields
+ // and immediately recombining them is an identity operation with
several divisions.
+ _data[old_size + row] = static_cast<TimeValue::TimeType>(micros);
Review Comment:
[P2] Reject Parquet TIME values outside the time-of-day domain
The Parquet [TIME
contract](https://github.com/apache/parquet-format/blob/master/LogicalTypes.md#time)
represents `TIME(MILLIS/MICROS/NANOS)` as units after midnight, but this
assignment publishes negative and at/above-one-day carriers as Doris signed
durations; the new unit test even expects `-1000001` to become
`-00:00:01.000001`. The same consumer is used by plain/sparse decode and typed
dictionaries, so malformed external input is silently returned as a different
logical value. Please validate the raw carrier against `[0, 24h)` in its
declared unit before rescaling, route nullable non-strict failures through
`ParquetMaterializationState` (including deferred dictionary-ID failures), and
add lower/upper boundary coverage for every unit.
##########
be/src/format_v2/parquet/parquet_statistics.cpp:
##########
@@ -762,435 +500,410 @@ void accumulate_zonemap_stats(const ZoneMapEvalContext&
ctx, ParquetPruningStats
} // namespace
+bool can_use_parquet_page_index(const format::FileScanRequest& request,
+ const RuntimeState* runtime_state) {
+ return config::enable_parquet_page_index &&
has_expr_zonemap_filter(request, runtime_state);
+}
+
std::shared_ptr<segment_v2::ZoneMap> ParquetStatisticsUtils::MakeZoneMap(
const ParquetColumnStatistics& statistics) {
return make_zonemap_from_statistics(statistics);
}
ParquetColumnStatistics ParquetStatisticsUtils::TransformColumnStatistics(
- const ParquetColumnSchema& column_schema,
- const std::shared_ptr<::parquet::Statistics>& statistics, const
cctz::time_zone* timezone) {
+ const ParquetColumnSchema& column_schema, const tparquet::Statistics*
statistics,
+ int64_t column_value_count, const cctz::time_zone* timezone) {
ParquetColumnStatistics result;
- if (statistics == nullptr) {
+ if (statistics == nullptr || column_value_count < 0) {
return result;
}
- result.has_null = !statistics->HasNullCount() || statistics->null_count()
> 0;
- result.has_not_null = statistics->num_values() > 0 ||
statistics->HasMinMax();
- result.has_null_count = statistics->HasNullCount();
- if (!result.has_not_null || !statistics->HasMinMax()) {
+ if (statistics->__isset.null_count && statistics->null_count >
column_value_count) {
+ // An impossible null count makes all derived min/max and all-null
flags untrustworthy;
+ // disable pruning instead of turning corrupt footer metadata into
false negatives.
return result;
}
- DORIS_CHECK(column_schema.type != nullptr);
- switch (statistics->physical_type()) {
- case ::parquet::Type::BOOLEAN:
- result.has_min_max = set_decoded_min_max<::parquet::BooleanType>(
- statistics, column_schema, DecodedValueKind::BOOL, &result,
timezone);
- return result;
- case ::parquet::Type::INT32:
- result.has_min_max = set_decoded_min_max<::parquet::Int32Type>(
- statistics, column_schema,
decoded_value_kind(column_schema.type_descriptor),
- &result, timezone);
- return result;
- case ::parquet::Type::INT64:
- result.has_min_max = set_decoded_min_max<::parquet::Int64Type>(
- statistics, column_schema,
decoded_value_kind(column_schema.type_descriptor),
- &result, timezone);
- return result;
- case ::parquet::Type::FLOAT:
- result.has_min_max = set_decoded_min_max<::parquet::FloatType>(
- statistics, column_schema, DecodedValueKind::FLOAT, &result,
timezone);
- return result;
- case ::parquet::Type::DOUBLE:
- result.has_min_max = set_decoded_min_max<::parquet::DoubleType>(
- statistics, column_schema, DecodedValueKind::DOUBLE, &result,
timezone);
- return result;
- case ::parquet::Type::BYTE_ARRAY:
- case ::parquet::Type::FIXED_LEN_BYTE_ARRAY:
- result.has_min_max = set_string_min_max(statistics, column_schema,
&result, timezone);
- return result;
- default:
- return result;
+ const bool has_null_count = statistics->__isset.null_count &&
statistics->null_count >= 0;
+ const int64_t null_count = has_null_count ? statistics->null_count : 0;
+ const bool has_not_null = has_null_count ? column_value_count > null_count
: true;
+ const std::string* min_value = statistics->__isset.min_value
+ ? &statistics->min_value
+ : (statistics->__isset.min ?
&statistics->min : nullptr);
+ const std::string* max_value = statistics->__isset.max_value
+ ? &statistics->max_value
+ : (statistics->__isset.max ?
&statistics->max : nullptr);
+
+ tparquet::ColumnIndex index;
+ index.__set_null_pages({!has_not_null});
+ index.__set_null_counts({null_count});
+ if (min_value != nullptr && max_value != nullptr) {
+ index.__set_min_values({*min_value});
+ index.__set_max_values({*max_value});
+ }
+ // Footer statistics and page indexes share the same little-endian
physical encoding. Reusing
+ // one decoder keeps native row-group and page pruning identical for
logical types and NaNs.
+ if (!build_native_page_statistics(index, column_schema, 0, &result,
timezone)) {
+ return {};
+ }
+ if (!has_null_count) {
+ result.has_null_count = false;
+ result.has_null = true;
}
+ return result;
}
-bool ParquetStatisticsUtils::BloomFilterExcludes(const ParquetColumnSchema&
column_schema,
- int slot_index, const
VExprContextSPtrs& conjuncts,
- const ::parquet::BloomFilter&
bloom_filter) {
- return bloom_filter_excludes(column_schema, slot_index, conjuncts,
bloom_filter);
+bool ParquetStatisticsUtils::NativeBloomFilterExcludes(
+ const ParquetColumnSchema& column_schema, int slot_index,
+ const VExprContextSPtrs& conjuncts, const segment_v2::BloomFilter&
bloom_filter) {
+ if (!bloom_filter_supported(column_schema)) {
+ return false;
+ }
+ NativeParquetBloomFilterAdapter adapter(column_schema, bloom_filter);
+ BloomFilterEvalContext ctx;
+ ctx.slots.emplace(slot_index, BloomFilterEvalContext::SlotBloomFilter {
+ .data_type = column_schema.type,
+ .bloom_filter = &adapter,
+ });
+ return VExprContext::evaluate_bloom_filter(conjuncts, ctx) ==
ZoneMapFilterResult::kNoMatch;
}
namespace {
-ParquetRowGroupPruneReason dictionary_prune_reason(
- const ::parquet::RowGroupMetaData& row_group,
::parquet::ParquetFileReader* file_reader,
- int row_group_idx, const
std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
- const format::FileScanRequest& request) {
- const auto conjuncts_by_slot = collect_conjuncts_by_single_slot(
- request.conjuncts, expr_zonemap::single_slot_dictionary_index);
- for (const auto& [slot_index, conjuncts] : conjuncts_by_slot) {
- const auto file_column_id = file_column_id_by_block_position(request,
slot_index);
- if (!file_column_id.has_value()) {
- continue;
- }
- const auto* column_schema = resolve_local_leaf_schema(file_schema,
*file_column_id);
- if (column_schema == nullptr || column_schema->type == nullptr) {
- continue;
- }
- DCHECK_LT(column_schema->leaf_column_id, row_group.num_columns());
- auto column_chunk =
row_group.ColumnChunk(column_schema->leaf_column_id);
- if (column_chunk == nullptr ||
- !supports_dictionary_pruning(*column_schema, *column_chunk) ||
- !is_dictionary_encoded_chunk(*column_chunk)) {
- continue;
- }
-
- ParquetDictionaryWords dict_words;
- if (!read_dictionary_words(file_reader, row_group_idx,
column_schema->leaf_column_id,
- *column_schema, &dict_words)) {
- continue;
- }
- DictionaryEvalContext ctx;
- ctx.slots.emplace(slot_index, DictionaryEvalContext::SlotDictionary {
- .data_type = column_schema->type,
- .values =
dictionary_fields_from_words(dict_words),
- });
- if (VExprContext::evaluate_dictionary_filter(conjuncts, ctx) ==
- ZoneMapFilterResult::kNoMatch) {
- return ParquetRowGroupPruneReason::DICTIONARY;
+void collect_filtered_leaf_ids(const ParquetColumnSchema& column_schema,
+ const format::LocalColumnIndex* projection,
+ std::set<int>* leaf_column_ids) {
+ if (column_schema.kind == ParquetColumnSchemaKind::PRIMITIVE) {
+ if (column_schema.leaf_column_id >= 0) {
+ leaf_column_ids->insert(column_schema.leaf_column_id);
}
+ return;
}
- return ParquetRowGroupPruneReason::NONE;
-}
-
-ParquetRowGroupPruneReason bloom_filter_prune_reason(
- int row_group_idx, const
std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
- const format::FileScanRequest& request, RowGroupBloomFilterCache*
bloom_filter_cache,
- ParquetPruningStats* pruning_stats) {
- if (bloom_filter_cache == nullptr) {
- return ParquetRowGroupPruneReason::NONE;
- }
- const auto conjuncts_by_slot = collect_conjuncts_by_single_slot(
- request.conjuncts, expr_zonemap::single_slot_bloom_filter_index);
- for (const auto& [slot_index, conjuncts] : conjuncts_by_slot) {
- const auto file_column_id = file_column_id_by_block_position(request,
slot_index);
- if (!file_column_id.has_value()) {
- continue;
- }
- const auto* column_schema = resolve_local_leaf_schema(file_schema,
*file_column_id);
- if (column_schema == nullptr || column_schema->type == nullptr ||
- !bloom_filter_supported(*column_schema)) {
- continue;
- }
- auto* bloom_filter = bloom_filter_cache->get(row_group_idx,
column_schema->leaf_column_id,
- pruning_stats);
- if (bloom_filter == nullptr) {
+ for (const auto& child_schema : column_schema.children) {
+ if (!format::is_child_projected(projection, child_schema->local_id)) {
continue;
}
- if (ParquetStatisticsUtils::BloomFilterExcludes(*column_schema,
slot_index, conjuncts,
- *bloom_filter)) {
- return ParquetRowGroupPruneReason::BLOOM_FILTER;
- }
+ collect_filtered_leaf_ids(*child_schema,
+ format::find_child_projection(projection,
child_schema->local_id),
+ leaf_column_ids);
}
- return ParquetRowGroupPruneReason::NONE;
}
-void init_bloom_filter_cache(::parquet::ParquetFileReader* file_reader, bool
enable_bloom_filter,
- RowGroupBloomFilterCache* bloom_filter_cache) {
- DORIS_CHECK(bloom_filter_cache != nullptr);
- if (!enable_bloom_filter || file_reader == nullptr) {
- return;
- }
- try {
- bloom_filter_cache->bloom_filter_reader =
&file_reader->GetBloomFilterReader();
- } catch (const ::parquet::ParquetException&) {
- bloom_filter_cache->bloom_filter_reader = nullptr;
- } catch (const std::exception&) {
- bloom_filter_cache->bloom_filter_reader = nullptr;
- }
+bool native_metadata_predicate_is_type_safe(const ParquetColumnSchema&
column_schema) {
+ DORIS_CHECK(column_schema.type != nullptr);
+ // Raw VARBINARY file slots may feed table-side STRING casts. Footer/page
metadata is still in
+ // the pre-cast domain, so using it for a rewritten table predicate can
cause false negatives.
+ return remove_nullable(column_schema.type)->get_primitive_type() !=
TYPE_VARBINARY;
}
-bool check_statistics(const ::parquet::RowGroupMetaData& row_group,
- const std::vector<std::unique_ptr<ParquetColumnSchema>>&
file_schema,
- const format::FileScanRequest& request,
ParquetPruningStats* pruning_stats,
- const cctz::time_zone* timezone) {
+bool check_native_statistics(const tparquet::RowGroup& row_group,
+ const
std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
+ const format::FileScanRequest& request,
+ ParquetPruningStats* pruning_stats, const
cctz::time_zone* timezone) {
const auto slot_indexes =
collect_expr_zonemap_slot_indexes(request.conjuncts);
if (slot_indexes.empty()) {
return false;
}
-
ZoneMapEvalContext ctx;
for (const int slot_index : slot_indexes) {
const auto file_column_id = file_column_id_by_block_position(request,
slot_index);
if (!file_column_id.has_value()) {
continue;
}
const auto* column_schema = resolve_local_leaf_schema(file_schema,
*file_column_id);
- if (column_schema == nullptr || column_schema->type == nullptr) {
+ if (column_schema == nullptr || column_schema->type == nullptr ||
+ !native_metadata_predicate_is_type_safe(*column_schema) ||
+ column_schema->leaf_column_id >=
static_cast<int>(row_group.columns.size())) {
continue;
}
-
+ const auto& chunk = row_group.columns[column_schema->leaf_column_id];
std::shared_ptr<segment_v2::ZoneMap> zone_map;
- DCHECK_LT(column_schema->leaf_column_id, row_group.num_columns());
- auto column_chunk =
row_group.ColumnChunk(column_schema->leaf_column_id);
- if (column_chunk != nullptr) {
+ if (chunk.__isset.meta_data) {
+ const auto& column_metadata = chunk.meta_data;
+ const auto* statistics =
+ column_metadata.__isset.statistics ?
&column_metadata.statistics : nullptr;
+ if (statistics != nullptr &&
!detail::can_use_native_footer_min_max(
+
column_schema->type_descriptor, *statistics)) {
+ statistics = nullptr;
+ }
zone_map = ParquetStatisticsUtils::MakeZoneMap(
ParquetStatisticsUtils::TransformColumnStatistics(
- *column_schema, column_chunk->statistics(),
timezone));
+ *column_schema, statistics,
column_metadata.num_values, timezone));
}
add_slot_zonemap(&ctx, slot_index, column_schema->type,
std::move(zone_map));
}
-
const auto result =
VExprContext::evaluate_zonemap_filter(request.conjuncts, ctx);
accumulate_zonemap_stats(ctx, pruning_stats);
return result == ZoneMapFilterResult::kNoMatch;
}
-Status select_row_groups_by_metadata_impl(
- const ::parquet::FileMetaData& metadata, ::parquet::ParquetFileReader*
file_reader,
- const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
- const format::FileScanRequest& request, const std::vector<int>*
candidate_row_groups,
- std::vector<int>* selected_row_groups, bool enable_bloom_filter,
- ParquetPruningStats* pruning_stats, const cctz::time_zone* timezone,
- const RuntimeState* runtime_state) {
- int64_t row_group_filter_time_sink = 0;
- SCOPED_RAW_TIMER(pruning_stats == nullptr ? &row_group_filter_time_sink
- :
&pruning_stats->row_group_filter_time);
- if (selected_row_groups == nullptr) {
- return Status::InvalidArgument("selected_row_groups is null");
- }
- selected_row_groups->clear();
+bool is_native_dictionary_data_encoding(tparquet::Encoding::type encoding) {
+ return encoding == tparquet::Encoding::PLAIN_DICTIONARY ||
+ encoding == tparquet::Encoding::RLE_DICTIONARY;
+}
- const int num_row_groups = metadata.num_row_groups();
- const auto candidate_size = candidate_row_groups == nullptr
- ? static_cast<size_t>(num_row_groups)
- : candidate_row_groups->size();
- if (pruning_stats != nullptr) {
- // Scan-range ownership is decided before metadata pruning. Count only
row groups owned by
- // this split so a file divided into multiple splits does not report
the full-file total and
- // out-of-split groups once per split.
- pruning_stats->total_row_groups = cast_set<int64_t>(candidate_size);
+bool is_native_level_encoding(tparquet::Encoding::type encoding) {
+ return encoding == tparquet::Encoding::RLE || encoding ==
tparquet::Encoding::BIT_PACKED;
+}
+
+bool is_native_dictionary_encoded_chunk(const tparquet::ColumnMetaData&
metadata) {
+ if (!metadata.__isset.dictionary_page_offset ||
metadata.dictionary_page_offset < 0) {
+ return false;
}
- selected_row_groups->reserve(candidate_size);
- RowGroupBloomFilterCache bloom_filter_cache;
- init_bloom_filter_cache(file_reader, enable_bloom_filter,
&bloom_filter_cache);
- for (size_t candidate_idx = 0; candidate_idx < candidate_size;
++candidate_idx) {
- const int row_group_idx = candidate_row_groups == nullptr
- ? static_cast<int>(candidate_idx)
- :
(*candidate_row_groups)[candidate_idx];
- DORIS_CHECK(row_group_idx >= 0);
- DORIS_CHECK(row_group_idx < num_row_groups);
- auto row_group = metadata.RowGroup(row_group_idx);
- if (row_group == nullptr) {
- selected_row_groups->push_back(row_group_idx);
- continue;
+ if (metadata.__isset.encoding_stats && !metadata.encoding_stats.empty()) {
+ bool has_dictionary_data_page = false;
+ for (const auto& encoding_stat : metadata.encoding_stats) {
+ if ((encoding_stat.page_type != tparquet::PageType::DATA_PAGE &&
+ encoding_stat.page_type != tparquet::PageType::DATA_PAGE_V2)
||
+ encoding_stat.count <= 0) {
+ continue;
+ }
+ if (!is_native_dictionary_data_encoding(encoding_stat.encoding)) {
+ return false;
+ }
+ has_dictionary_data_page = true;
}
- ParquetRowGroupPruneReason prune_reason =
ParquetRowGroupPruneReason::NONE;
- if (has_expr_zonemap_filter(request, runtime_state) &&
- check_statistics(*row_group, file_schema, request, pruning_stats,
timezone)) {
- prune_reason = ParquetRowGroupPruneReason::STATISTICS;
+ return has_dictionary_data_page;
+ }
+ bool has_dictionary_encoding = false;
+ for (const auto encoding : metadata.encodings) {
+ if (is_native_dictionary_data_encoding(encoding)) {
+ has_dictionary_encoding = true;
+ } else if (!is_native_level_encoding(encoding)) {
+ return false;
}
+ }
+ return has_dictionary_encoding;
+}
- if (prune_reason == ParquetRowGroupPruneReason::NONE) {
- prune_reason = dictionary_prune_reason(*row_group, file_reader,
row_group_idx,
- file_schema, request);
- if (prune_reason == ParquetRowGroupPruneReason::NONE) {
- prune_reason = bloom_filter_prune_reason(row_group_idx,
file_schema, request,
- &bloom_filter_cache,
pruning_stats);
- }
+const format::LocalColumnIndex* find_request_projection(const
format::FileScanRequest& request,
+ format::LocalColumnId
file_column_id) {
+ for (const auto& projection : request.predicate_columns) {
+ if (projection.local_id() == file_column_id.value()) {
+ return &projection;
}
-
- if (prune_reason != ParquetRowGroupPruneReason::NONE) {
- if (pruning_stats != nullptr) {
- pruning_stats->filtered_group_rows += row_group->num_rows();
- if (prune_reason == ParquetRowGroupPruneReason::STATISTICS) {
- ++pruning_stats->filtered_row_groups_by_statistics;
- } else if (prune_reason ==
ParquetRowGroupPruneReason::DICTIONARY) {
- ++pruning_stats->filtered_row_groups_by_dictionary;
- } else if (prune_reason ==
ParquetRowGroupPruneReason::BLOOM_FILTER) {
- ++pruning_stats->filtered_row_groups_by_bloom_filter;
- }
- }
- continue;
+ }
+ for (const auto& projection : request.non_predicate_columns) {
+ if (projection.local_id() == file_column_id.value()) {
+ return &projection;
}
- selected_row_groups->push_back(row_group_idx);
}
- return Status::OK();
+ return nullptr;
}
-} // namespace
-
-Status select_row_groups_by_metadata(
- const ::parquet::FileMetaData& metadata, ::parquet::ParquetFileReader*
file_reader,
+ParquetRowGroupPruneReason native_dictionary_prune_reason(
+ const tparquet::RowGroup& row_group, int row_group_idx,
const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
- const format::FileScanRequest& request, const std::vector<int>*
candidate_row_groups,
- std::vector<int>* selected_row_groups, bool enable_bloom_filter,
- ParquetPruningStats* pruning_stats, const cctz::time_zone* timezone,
- const RuntimeState* runtime_state) {
- return select_row_groups_by_metadata_impl(
- metadata, file_reader, file_schema, request, candidate_row_groups,
selected_row_groups,
- enable_bloom_filter, pruning_stats, timezone, runtime_state);
-}
-
-namespace {
-
-template <typename ParquetDType>
-bool set_page_decoded_min_max(const std::shared_ptr<::parquet::ColumnIndex>&
column_index,
- const ParquetColumnSchema& column_schema, size_t
page_idx,
- DecodedValueKind value_kind,
ParquetColumnStatistics* page_statistics,
- const cctz::time_zone* timezone) {
- const auto typed_index =
-
std::static_pointer_cast<::parquet::TypedColumnIndex<ParquetDType>>(column_index);
- if (page_idx >= typed_index->min_values().size() ||
- page_idx >= typed_index->max_values().size()) {
- return false;
+ const format::FileScanRequest& request, const cctz::time_zone*
timezone,
+ ParquetFileContext* file_context, const ParquetColumnReaderProfile&
column_reader_profile) {
+ if (file_context == nullptr || file_context->native_metadata == nullptr) {
+ return ParquetRowGroupPruneReason::NONE;
}
- const auto& min_value = typed_index->min_values()[page_idx];
- const auto& max_value = typed_index->max_values()[page_idx];
- if constexpr (std::is_same_v<ParquetDType, ::parquet::Int64Type>) {
- if (!timestamp_min_max_is_safe(column_schema, min_value, max_value,
timezone)) {
- return false;
+ const auto conjuncts_by_slot = collect_conjuncts_by_single_slot(
+ request.conjuncts, expr_zonemap::single_slot_dictionary_index);
+ for (const auto& [slot_index, conjuncts] : conjuncts_by_slot) {
+ const auto file_column_id = file_column_id_by_block_position(request,
slot_index);
+ if (!file_column_id.has_value()) {
+ continue;
+ }
+ const auto* column_schema = resolve_local_leaf_schema(file_schema,
*file_column_id);
+ const auto* projection = find_request_projection(request,
*file_column_id);
+ if (column_schema == nullptr || projection == nullptr ||
column_schema->type == nullptr ||
+ !column_schema->type_descriptor.is_string_like ||
+ column_schema->leaf_column_id >=
static_cast<int>(row_group.columns.size())) {
+ continue;
+ }
+ if (!native_metadata_predicate_is_type_safe(*column_schema)) {
+ // The file-local VARBINARY may feed a table-side STRING cast.
Pruning before that cast
+ // can compare different Field kinds and incorrectly discard a
matching row group.
+ continue;
+ }
+ const auto& chunk = row_group.columns[column_schema->leaf_column_id];
+ if (!chunk.__isset.meta_data ||
+ (chunk.meta_data.type != tparquet::Type::BYTE_ARRAY &&
+ chunk.meta_data.type != tparquet::Type::FIXED_LEN_BYTE_ARRAY) ||
+ !is_native_dictionary_encoded_chunk(chunk.meta_data)) {
+ continue;
+ }
+ std::unique_ptr<ParquetColumnReader> reader;
+ const std::vector<RowRange> ranges {{0, row_group.num_rows}};
Review Comment:
[P1] Skip empty row groups before dictionary probing
`NativeParquetMetadata::init_schema()` accepts `num_rows == 0`, and
`plan_parquet_row_groups()` runs metadata pruning before
`build_native_row_group_read_plans()` performs its intended empty-group skip.
For a string predicate on a dictionary-marked empty chunk, this creates `{0,
0}` and enters `NativeColumnReader::init()`, whose
`DORIS_CHECK(row_group.num_rows > 0)` (and positive-range check) terminates the
BE. The existing native reader even documents that Parquet may emit empty row
groups. Please discard zero-row groups before all metadata probes (or return
`NONE` here without constructing a reader), and add a dictionary-marked
zero-row planning test; inconsistent external metadata should return
`Status::Corruption`, never abort the process.
##########
be/src/format_v2/delimited_text/delimited_text_reader.cpp:
##########
@@ -200,7 +203,7 @@ void DelimitedTextReader::_init_profile() {
_text_profile.rows_read_before_filter = ADD_CHILD_COUNTER_WITH_LEVEL(
_profile, "RowsReadBeforeFilter", TUnit::UNIT,
DELIMITED_TEXT_PROFILE, 1);
_text_profile.rows_filtered_by_conjunct = ADD_CHILD_COUNTER_WITH_LEVEL(
- _profile, "RowsFilteredByConjunct", TUnit::UNIT,
DELIMITED_TEXT_PROFILE, 1);
+ _profile, "RowsFilteredByConjunct", TUnit::UNIT,
file_scan_profile::FILE_READER, 1);
Review Comment:
[P2] Keep this delimited counter format-specific
`RuntimeProfile` keys counters globally by name, and Parquet registers
`RowsFilteredByConjunct` below `ParquetReader`. Because one `TableReader` can
switch formats across splits while reusing the scanner profile, this move makes
ownership initialization-order dependent: delimited-first leaves Parquet
updates under `FileReader`, while Parquet-first attributes CSV/TEXT updates
under `ParquetReader`. It also removes the metric from `DelimitedTextReader`
even in single-format scans. Please retain a uniquely named delimited child (or
publish a separate common aggregate) and cover both mixed initialization orders
plus each format subtree.
##########
be/src/format_v2/parquet/parquet_profile.cpp:
##########
@@ -69,26 +80,41 @@ void ParquetProfile::init(RuntimeProfile* profile) {
parquet_profile, 1);
reader_select_rows = ADD_CHILD_COUNTER_WITH_LEVEL(profile,
"ReaderSelectRows", TUnit::UNIT,
parquet_profile, 1);
- arrow_read_records_time =
- ADD_CHILD_TIMER_WITH_LEVEL(profile, "ArrowReadRecordsTime",
parquet_profile, 1);
- arrow_skip_records_time =
- ADD_CHILD_TIMER_WITH_LEVEL(profile, "ArrowSkipRecordsTime",
parquet_profile, 1);
+ level_only_read_time =
+ ADD_CHILD_TIMER_WITH_LEVEL(profile, "LevelOnlyReadTime",
parquet_profile, 1);
+ level_only_skip_time =
+ ADD_CHILD_TIMER_WITH_LEVEL(profile, "LevelOnlySkipTime",
parquet_profile, 1);
materialization_time =
ADD_CHILD_TIMER_WITH_LEVEL(profile, "MaterializationTime",
parquet_profile, 1);
+ hybrid_selection_batches = ADD_CHILD_COUNTER_WITH_LEVEL(profile,
"HybridSelectionBatches",
+ TUnit::UNIT,
parquet_profile, 1);
+ hybrid_selection_ranges = ADD_CHILD_COUNTER_WITH_LEVEL(profile,
"HybridSelectionRanges",
+ TUnit::UNIT,
parquet_profile, 1);
+ hybrid_selection_null_fallback_batches = ADD_CHILD_COUNTER_WITH_LEVEL(
+ profile, "HybridSelectionNullFallbackBatches", TUnit::UNIT,
parquet_profile, 1);
+ native_read_calls = ADD_CHILD_COUNTER_WITH_LEVEL(profile,
"NativeReadCalls", TUnit::UNIT,
+ parquet_profile, 1);
+ native_page_fragments = ADD_CHILD_COUNTER_WITH_LEVEL(profile,
"NativePageFragments",
+ TUnit::UNIT,
parquet_profile, 1);
+ page_crossing_batches = ADD_CHILD_COUNTER_WITH_LEVEL(profile,
"PageCrossingBatches",
+ TUnit::UNIT,
parquet_profile, 1);
+ nested_batches =
+ ADD_CHILD_COUNTER_WITH_LEVEL(profile, "NestedBatches",
TUnit::UNIT, parquet_profile, 1);
lazy_read_filtered_rows = ADD_CHILD_COUNTER_WITH_LEVEL(profile,
"FilteredRowsByLazyRead",
TUnit::UNIT,
parquet_profile, 1);
filtered_bytes = ADD_CHILD_COUNTER_WITH_LEVEL(profile, "FilteredBytes",
TUnit::BYTES,
- parquet_profile, 1);
+
file_scan_profile::FILE_READER, 1);
Review Comment:
[P2] Keep these Parquet metrics in the ParquetReader subtree
`FilteredBytes` (and `FileNum` below) are populated only by Parquet
planning/open paths, but parenting them under `FileReader` makes them siblings
of `ParquetReader`. That contradicts the documented format-subtree ownership
and means consumers that extract `ParquetReader` silently lose the avoided-byte
and opened-file metrics. The hierarchy test never asserts either displaced
counter's parent, while a separate profile test checks `FileNum` only through
the flat name lookup. Please parent both under `parquet_profile` (or preserve
format-owned children alongside separately named common totals) and assert
their hierarchy explicitly.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]