This is an automated email from the ASF dual-hosted git repository.
Gabriel39 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new b80247f1200 [fix](file) Resolve remaining FileScannerV2 audit issues
(#65931)
b80247f1200 is described below
commit b80247f1200087aa4985ad7c67ae2cbed82b33ce
Author: Gabriel <[email protected]>
AuthorDate: Thu Jul 23 21:47:46 2026 +0800
[fix](file) Resolve remaining FileScannerV2 audit issues (#65931)
### What problem does this PR solve?
This follow-up audits the 10 unresolved review threads on
https://github.com/apache/doris/pull/65674 and every child of
DORIS-27038.
- Three unresolved threads were already fixed on master.
- This PR fixes the remaining seven review findings.
### What is changed?
- Harden Parquet delta geometry, page allocation validation, schema
ambiguity handling, dictionary reuse, decompression scratch lifetime,
and one-child MAP_KEY_VALUE SET parsing.
- Validate delete expression result ownership before erasing temporary
columns.
- Scope Iceberg row-lineage virtual columns to Iceberg readers.
- Propagate Remote Doris Flight timeout and cancellation.
- Reject NULL JDBC special-type casts for non-nullable targets.
- Aggregate hybrid Paimon/Hudi condition-cache hit counters.
- Remove redundant JSON-line copies and Hive key allocations.
The scanner V1 path under `be/src/format` is intentionally untouched.
### Verification
- BE ASAN unit tests: 172 tests from 12 related suites passed.
- clang-format 16 dry-run passed for every changed C++ file.
- `git diff --check` passed.
- Confirmed no diff under `be/src/format`.
---------
Signed-off-by: Gabriel <[email protected]>
---
be/src/format_v2/column_mapper.cpp | 9 +-
be/src/format_v2/column_mapper.h | 1 +
be/src/format_v2/jni/jdbc_reader.cpp | 15 +
be/src/format_v2/jni/jdbc_reader.h | 2 +
be/src/format_v2/json/json_reader.cpp | 54 ++-
be/src/format_v2/json/json_reader.h | 44 ++-
be/src/format_v2/parquet/native_schema_desc.cpp | 16 +-
be/src/format_v2/parquet/native_schema_node.cpp | 35 +-
be/src/format_v2/parquet/parquet_scan.cpp | 23 +-
be/src/format_v2/parquet/parquet_scan.h | 2 +
.../parquet/reader/native/column_chunk_reader.cpp | 162 ++++++++-
.../parquet/reader/native/column_chunk_reader.h | 32 +-
.../parquet/reader/native/column_reader.cpp | 16 +-
.../parquet/reader/native/column_reader.h | 6 +
.../parquet/reader/native/delta_bit_pack_decoder.h | 6 +
.../parquet/reader/native_column_reader.cpp | 19 +-
.../parquet/reader/native_column_reader.h | 1 +
be/src/format_v2/table/hudi_reader.cpp | 7 +
be/src/format_v2/table/hudi_reader.h | 5 +
.../iceberg_position_delete_sys_table_reader.cpp | 1 +
be/src/format_v2/table/iceberg_reader.h | 1 +
be/src/format_v2/table/paimon_reader.cpp | 7 +
be/src/format_v2/table/paimon_reader.h | 5 +
be/src/format_v2/table/remote_doris_reader.cpp | 185 +++++++++-
be/src/format_v2/table/remote_doris_reader.h | 2 +
be/src/format_v2/table_reader.h | 3 +-
be/test/format_v2/column_mapper_test.cpp | 28 +-
be/test/format_v2/jni/jdbc_reader_test.cpp | 49 +++
be/test/format_v2/json/json_reader_test.cpp | 24 +-
be/test/format_v2/parquet/native_decoder_test.cpp | 381 ++++++++++++++++++++-
be/test/format_v2/parquet/parquet_scan_test.cpp | 50 +++
be/test/format_v2/parquet/parquet_schema_test.cpp | 112 ++++++
be/test/format_v2/table/hudi_reader_test.cpp | 7 +
be/test/format_v2/table/paimon_reader_test.cpp | 7 +
.../format_v2/table/remote_doris_reader_test.cpp | 269 +++++++++++++++
35 files changed, 1511 insertions(+), 75 deletions(-)
diff --git a/be/src/format_v2/column_mapper.cpp
b/be/src/format_v2/column_mapper.cpp
index 8ac3060bac4..ed0dcda9739 100644
--- a/be/src/format_v2/column_mapper.cpp
+++ b/be/src/format_v2/column_mapper.cpp
@@ -395,7 +395,7 @@ std::string TableColumnMapperOptions::debug_string() const {
std::ostringstream out;
out << "TableColumnMapperOptions{mode=" << mapping_mode_to_string(mode)
<< ", allow_idless_complex_wrapper_projection=" <<
allow_idless_complex_wrapper_projection
- << "}";
+ << ", enable_row_lineage_virtual_columns=" <<
enable_row_lineage_virtual_columns << "}";
return out.str();
}
@@ -2031,7 +2031,12 @@ Status
TableColumnMapper::_create_mapping_for_column(const ColumnDefinition& tab
mapping->global_index = global_index;
mapping->table_column_name = table_column.name;
mapping->table_type = table_column.type;
- const auto row_lineage_type =
row_lineage_virtual_column_type(table_column, _options.mode);
+ // Row-lineage names are Iceberg metadata contracts, not reserved names in
generic Hive,
+ // Hudi, or Paimon schemas. Only the Iceberg reader may opt into virtual
synthesis.
+ const auto row_lineage_type =
+ _options.enable_row_lineage_virtual_columns
+ ? row_lineage_virtual_column_type(table_column,
_options.mode)
+ : TableVirtualColumnType::INVALID;
if (const auto* partition_value = find_partition_value(table_column,
_partition_values);
table_column.is_partition_key && partition_value != nullptr) {
// Partition values are split constants and must take precedence over
defaults.
diff --git a/be/src/format_v2/column_mapper.h b/be/src/format_v2/column_mapper.h
index 1ee306a8774..0860b687e4a 100644
--- a/be/src/format_v2/column_mapper.h
+++ b/be/src/format_v2/column_mapper.h
@@ -167,6 +167,7 @@ struct ColumnMapping {
struct TableColumnMapperOptions {
TableColumnMappingMode mode = TableColumnMappingMode::BY_FIELD_ID;
bool allow_idless_complex_wrapper_projection = false;
+ bool enable_row_lineage_virtual_columns = false;
std::string debug_string() const;
};
diff --git a/be/src/format_v2/jni/jdbc_reader.cpp
b/be/src/format_v2/jni/jdbc_reader.cpp
index 7d28134db4d..41d3a5282b6 100644
--- a/be/src/format_v2/jni/jdbc_reader.cpp
+++ b/be/src/format_v2/jni/jdbc_reader.cpp
@@ -34,6 +34,20 @@
namespace doris::format::jdbc {
+Status validate_non_nullable_special_type_result(const IColumn& result, size_t
rows) {
+ const auto* nullable = check_and_get_column<ColumnNullable>(&result);
+ if (UNLIKELY(nullable == nullptr)) {
+ return Status::InternalError("JDBC special-type CAST did not return a
nullable column");
+ }
+ if (UNLIKELY(nullable->has_null(0, rows))) {
+ // CAST NULL represents invalid source data; stripping the null map
would turn it into a
+ // valid-looking default for a NOT NULL destination.
+ return Status::DataQualityError(
+ "JDBC special-type CAST produced NULL for a non-nullable
column");
+ }
+ return Status::OK();
+}
+
std::string JdbcJniReader::connector_class() const {
return "org/apache/doris/jdbc/JdbcJniScanner";
}
@@ -184,6 +198,7 @@ Status JdbcJniReader::_cast_string_to_special_type(const
format::JniTableReader:
if (target_type->is_nullable()) {
output_block->replace_by_position(column.output_index, result_column);
} else {
+
RETURN_IF_ERROR(validate_non_nullable_special_type_result(*result_column,
rows));
const auto* nullable_column = assert_cast<const
ColumnNullable*>(result_column.get());
output_block->replace_by_position(column.output_index,
nullable_column->get_nested_column_ptr());
diff --git a/be/src/format_v2/jni/jdbc_reader.h
b/be/src/format_v2/jni/jdbc_reader.h
index 91a5878cb46..e10f8ad9790 100644
--- a/be/src/format_v2/jni/jdbc_reader.h
+++ b/be/src/format_v2/jni/jdbc_reader.h
@@ -29,6 +29,8 @@
namespace doris::format::jdbc {
+Status validate_non_nullable_special_type_result(const IColumn& result, size_t
rows);
+
class JdbcJniReader final : public format::JniTableReader {
public:
~JdbcJniReader() override = default;
diff --git a/be/src/format_v2/json/json_reader.cpp
b/be/src/format_v2/json/json_reader.cpp
index 1f42a2cd2d5..caf4205fa61 100644
--- a/be/src/format_v2/json/json_reader.cpp
+++ b/be/src/format_v2/json/json_reader.cpp
@@ -270,10 +270,16 @@ Status JsonReader::open(std::shared_ptr<FileScanRequest>
request) {
DORIS_CHECK(_request != nullptr);
RETURN_IF_ERROR(_build_requested_columns(*_request, &_requested_columns));
_slot_name_to_index.clear();
+ _hive_slot_name_to_index.clear();
_slot_name_to_index.reserve(_requested_columns.size());
+ _hive_slot_name_to_index.reserve(_requested_columns.size());
for (size_t idx = 0; idx < _requested_columns.size(); ++idx) {
auto name = _requested_columns[idx].slot_desc->col_name();
- _slot_name_to_index.emplace(_is_hive_table ? lower_key(name) : name,
idx);
+ if (_is_hive_table) {
+ _hive_slot_name_to_index.emplace(std::move(name), idx);
+ } else {
+ _slot_name_to_index.emplace(std::move(name), idx);
+ }
}
_previous_positions.clear();
_reader_range = _json_range();
@@ -532,7 +538,9 @@ Status JsonReader::_read_one_document(size_t* size, bool*
eof) {
if (*eof) {
return Status::OK();
}
- _document_buffer.assign(reinterpret_cast<const char*>(line), *size);
+ // The line reader owns this span until the next read, and parsing
copies it immediately
+ // into the padded buffer. Borrowing it avoids a redundant line-sized
string copy.
+ _document_view = std::string_view(reinterpret_cast<const char*>(line),
*size);
return Status::OK();
}
// Non-line mode treats the split as one JSON document. This supports a
single object or an
@@ -558,6 +566,7 @@ Status JsonReader::_read_one_document(size_t* size, bool*
eof) {
Slice result(_document_buffer.data(), _document_buffer.size());
RETURN_IF_ERROR(_physical_file_reader->read_at(_current_offset, result,
size, _io_ctx.get()));
_document_buffer.resize(*size);
+ _document_view = _document_buffer;
if (*size == 0) {
*eof = true;
}
@@ -572,6 +581,7 @@ Status JsonReader::_read_one_document_from_pipe(size_t*
read_size) {
DorisUniqueBufferPtr<uint8_t> file_buf;
RETURN_IF_ERROR(stream_load_pipe->read_one_message(&file_buf, read_size));
_document_buffer.assign(reinterpret_cast<const char*>(file_buf.get()),
*read_size);
+ _document_view = _document_buffer;
if (!stream_load_pipe->is_chunked_transfer()) {
return Status::OK();
}
@@ -586,6 +596,7 @@ Status JsonReader::_read_one_document_from_pipe(size_t*
read_size) {
_document_buffer.append(reinterpret_cast<const char*>(next_buf.get()),
next_size);
*read_size += next_size;
}
+ _document_view = _document_buffer;
return Status::OK();
}
@@ -594,10 +605,10 @@ Status JsonReader::_parse_next_json(size_t* size, bool*
eof) {
if (*eof || *size == 0) {
return Status::OK();
}
- if (*size >= 3 && static_cast<unsigned char>(_document_buffer[0]) == 0xEF
&&
- static_cast<unsigned char>(_document_buffer[1]) == 0xBB &&
- static_cast<unsigned char>(_document_buffer[2]) == 0xBF) {
- _document_buffer.erase(0, 3);
+ if (*size >= 3 && static_cast<unsigned char>(_document_view[0]) == 0xEF &&
+ static_cast<unsigned char>(_document_view[1]) == 0xBB &&
+ static_cast<unsigned char>(_document_view[2]) == 0xBF) {
+ _document_view.remove_prefix(3);
*size -= 3;
}
if (*size + simdjson::SIMDJSON_PADDING > _padded_size) {
@@ -606,7 +617,7 @@ Status JsonReader::_parse_next_json(size_t* size, bool*
eof) {
}
// Ondemand values reference the input buffer. Keep the padded bytes in a
member buffer until the
// current document is fully materialized.
- std::memcpy(_padding_buffer.data(), _document_buffer.data(), *size);
+ std::memcpy(_padding_buffer.data(), _document_view.data(), *size);
_original_doc_size = *size;
const auto error =
_json_parser->iterate(std::string_view(_padding_buffer.data(),
*size), _padded_size)
@@ -1120,32 +1131,39 @@ void JsonReader::_pop_back_last_inserted_value(Block*
block, size_t column_index
}
size_t JsonReader::_column_index(std::string_view key, size_t key_index) {
- std::string hive_key;
- std::string_view lookup_key = key;
- if (_is_hive_table) {
- hive_key = lower_key(key);
- lookup_key = hive_key;
- }
if (key_index < _previous_positions.size()) {
// Most JSON lines share field order. Reuse the previous line's
key-position mapping before
// falling back to the hash table lookup.
const auto previous = _previous_positions[key_index];
if (previous < _requested_columns.size()) {
const auto previous_name =
_requested_columns[previous].slot_desc->col_name();
- if ((_is_hive_table ? lower_key(previous_name) : previous_name) ==
lookup_key) {
+ if ((_is_hive_table && CaseInsensitiveStringEqual
{}(previous_name, key)) ||
+ (!_is_hive_table && previous_name == key)) {
return previous;
}
}
}
- const auto it = _slot_name_to_index.find(std::string(lookup_key));
- if (it == _slot_name_to_index.end()) {
+ // Transparent lookup keeps the common per-key path allocation-free; Hive
case folding is
+ // performed by the map's hash/equality without constructing a lower-cased
string.
+ const auto index = _is_hive_table ? ([&]() -> std::optional<size_t> {
+ const auto it = _hive_slot_name_to_index.find(key);
+ return it == _hive_slot_name_to_index.end() ? std::nullopt
+ :
std::optional<size_t>(it->second);
+ })()
+ : ([&]() -> std::optional<size_t> {
+ const auto it =
_slot_name_to_index.find(key);
+ return it ==
_slot_name_to_index.end()
+ ? std::nullopt
+ :
std::optional<size_t>(it->second);
+ })();
+ if (!index.has_value()) {
return static_cast<size_t>(-1);
}
if (key_index >= _previous_positions.size()) {
_previous_positions.resize(key_index + 1, static_cast<size_t>(-1));
}
- _previous_positions[key_index] = it->second;
- return it->second;
+ _previous_positions[key_index] = *index;
+ return *index;
}
bool JsonReader::_is_root_path_for_column(const RequestedColumn& column) const
{
diff --git a/be/src/format_v2/json/json_reader.h
b/be/src/format_v2/json/json_reader.h
index de3c084e681..c7346cb1d66 100644
--- a/be/src/format_v2/json/json_reader.h
+++ b/be/src/format_v2/json/json_reader.h
@@ -42,6 +42,40 @@ class IColumn;
namespace doris::format::json {
+struct TransparentStringHash {
+ using is_transparent = void;
+ size_t operator()(std::string_view value) const {
+ return std::hash<std::string_view> {}(value);
+ }
+};
+
+struct CaseInsensitiveStringHash {
+ using is_transparent = void;
+ size_t operator()(std::string_view value) const {
+ size_t hash = 1469598103934665603ULL;
+ for (const unsigned char c : value) {
+ const unsigned char folded = c >= 'A' && c <= 'Z' ? c + ('a' -
'A') : c;
+ hash = (hash ^ folded) * 1099511628211ULL;
+ }
+ return hash;
+ }
+};
+
+struct CaseInsensitiveStringEqual {
+ using is_transparent = void;
+ bool operator()(std::string_view lhs, std::string_view rhs) const {
+ if (lhs.size() != rhs.size()) return false;
+ for (size_t i = 0; i < lhs.size(); ++i) {
+ const unsigned char left =
+ lhs[i] >= 'A' && lhs[i] <= 'Z' ? lhs[i] + ('a' - 'A') :
lhs[i];
+ const unsigned char right =
+ rhs[i] >= 'A' && rhs[i] <= 'Z' ? rhs[i] + ('a' - 'A') :
rhs[i];
+ if (left != right) return false;
+ }
+ return true;
+ }
+};
+
// FileScannerV2 JSON reader.
//
// JSON files do not carry an embedded physical schema. The v2 table layer
still needs a
@@ -74,6 +108,10 @@ public:
Status get_block(Block* file_block, size_t* rows, bool* eof) override;
Status close() override;
+#ifdef BE_TEST
+ size_t TEST_document_buffer_size() const { return _document_buffer.size();
}
+#endif
+
private:
void _init_profile() override;
// A requested column keeps both identities:
@@ -135,7 +173,10 @@ private:
TFileCompressType::type _range_compress_type = TFileCompressType::UNKNOWN;
std::optional<TUniqueId> _stream_load_id;
std::vector<RequestedColumn> _requested_columns;
- std::unordered_map<std::string, size_t> _slot_name_to_index;
+ std::unordered_map<std::string, size_t, TransparentStringHash,
std::equal_to<>>
+ _slot_name_to_index;
+ std::unordered_map<std::string, size_t, CaseInsensitiveStringHash,
CaseInsensitiveStringEqual>
+ _hive_slot_name_to_index;
std::vector<size_t> _previous_positions;
RuntimeProfile::Counter* _total_time = nullptr;
@@ -178,6 +219,7 @@ private:
simdjson::ondemand::array _array;
simdjson::ondemand::array_iterator _array_iter;
std::string _document_buffer;
+ std::string_view _document_view;
std::string _padding_buffer;
size_t _original_doc_size = 0;
size_t _padded_size = 1024 * 1024 * 8 + simdjson::SIMDJSON_PADDING;
diff --git a/be/src/format_v2/parquet/native_schema_desc.cpp
b/be/src/format_v2/parquet/native_schema_desc.cpp
index 713707b84b8..9eba6d8feaa 100644
--- a/be/src/format_v2/parquet/native_schema_desc.cpp
+++ b/be/src/format_v2/parquet/native_schema_desc.cpp
@@ -207,10 +207,6 @@ static bool is_struct_list_node(const
tparquet::SchemaElement& schema,
return schema.name == "array" || schema.name == enclosing_list_name +
"_tuple";
}
-static bool has_logical_annotation(const tparquet::SchemaElement& schema) {
- return schema.__isset.logicalType || schema.__isset.converted_type;
-}
-
std::string NativeFieldSchema::debug_string() const {
std::stringstream ss;
ss << "NativeFieldSchema(name=" << name << ", R=" << repetition_level
@@ -625,9 +621,17 @@ Status NativeFieldDescriptor::parse_list_field(
if (num_children > 0) {
const bool structural_wrapper = is_struct_list_node(second_level,
first_level.name);
const auto& only_child = t_schemas[curr_pos + 2];
- if (num_children == 1 && !structural_wrapper &&
has_logical_annotation(second_level)) {
+ const bool single_key_set_wrapper =
+ second_level.__isset.converted_type &&
+ second_level.converted_type ==
tparquet::ConvertedType::MAP_KEY_VALUE &&
+ !is_repeated_node(only_child);
+ const bool nested_collection_annotation =
+ is_list_node(second_level) ||
+ (is_map_node(second_level) && !single_key_set_wrapper);
+ if (num_children == 1 && !structural_wrapper &&
nested_collection_annotation) {
// The repeated node is already the outer LIST element. Preserve
its own LIST/MAP
- // annotation, but do not interpret its REPEATED marker as another
outer array.
+ // annotation. A MAP_KEY_VALUE node with a repeated child is the
legacy uncontained
+ // MAP shape; only its direct, non-repeated child denotes the
enclosing SET key.
set_child_node_level(list_field, list_field->definition_level);
if (is_list_node(second_level)) {
RETURN_IF_ERROR(parse_list_field(t_schemas, curr_pos + 1,
list_child, true));
diff --git a/be/src/format_v2/parquet/native_schema_node.cpp
b/be/src/format_v2/parquet/native_schema_node.cpp
index 7539fb45efb..052df93e4b4 100644
--- a/be/src/format_v2/parquet/native_schema_node.cpp
+++ b/be/src/format_v2/parquet/native_schema_node.cpp
@@ -68,24 +68,39 @@ Status build_native_schema_node(const DataTypePtr&
projected_type,
return Status::Corruption("Parquet column {} is not a STRUCT",
file_schema.name);
}
const auto* struct_type = assert_cast<const
DataTypeStruct*>(type.get());
- std::map<std::string, const ParquetColumnSchema*> file_children;
- for (const auto& child : file_schema.children) {
- file_children.emplace(to_lower(child->name), child.get());
- }
auto node = std::make_shared<NativeStructSchemaNode>();
for (size_t i = 0; i < struct_type->get_elements().size(); ++i) {
const auto& table_name = struct_type->get_element_name(i);
- // Native metadata keeps writer casing. Match normalized names
while preserving the
- // original file name used to address the physical child reader.
- const auto child_it = file_children.find(to_lower(table_name));
- if (child_it == file_children.end()) {
+ const ParquetColumnSchema* file_child = nullptr;
+ for (const auto& child : file_schema.children) {
+ if (child->name == table_name) {
+ file_child = child.get();
+ break;
+ }
+ }
+ if (file_child == nullptr) {
+ for (const auto& child : file_schema.children) {
+ if (to_lower(child->name) != to_lower(table_name)) {
+ continue;
+ }
+ if (UNLIKELY(file_child != nullptr)) {
+ // Exact writer identity is authoritative; normalized
fallback is only safe
+ // when the requested field has one physical candidate.
+ return Status::Corruption(
+ "Parquet STRUCT {} has ambiguous
case-insensitive child name {}",
+ file_schema.name, table_name);
+ }
+ file_child = child.get();
+ }
+ }
+ if (file_child == nullptr) {
node->add_missing_child(table_name);
continue;
}
std::shared_ptr<NativeSchemaNode> child_node;
-
RETURN_IF_ERROR(build_native_schema_node(struct_type->get_element(i),
*child_it->second,
+
RETURN_IF_ERROR(build_native_schema_node(struct_type->get_element(i),
*file_child,
&child_node));
- node->add_child(table_name, child_it->second->name,
std::move(child_node));
+ node->add_child(table_name, file_child->name,
std::move(child_node));
}
*result = std::move(node);
return Status::OK();
diff --git a/be/src/format_v2/parquet/parquet_scan.cpp
b/be/src/format_v2/parquet/parquet_scan.cpp
index 47e7c4c3a68..b92a3c27054 100644
--- a/be/src/format_v2/parquet/parquet_scan.cpp
+++ b/be/src/format_v2/parquet/parquet_scan.cpp
@@ -652,11 +652,12 @@ Status execute_compact_delete_conjuncts(const
VExprContextSPtrs& delete_conjunct
*can_filter_all = false;
for (const auto& delete_conjunct : delete_conjuncts) {
DORIS_CHECK(delete_conjunct != nullptr);
+ const size_t original_columns = file_block->columns();
int result_column_id = -1;
RETURN_IF_ERROR(delete_conjunct->root()->execute(delete_conjunct.get(),
file_block,
&result_column_id));
- DORIS_CHECK(result_column_id >= 0 &&
- result_column_id <
static_cast<int>(file_block->columns()));
+ RETURN_IF_ERROR(detail::validate_ephemeral_expr_result_column(
+ original_columns, result_column_id, file_block->columns()));
const auto& delete_filter = assert_cast<const ColumnUInt8&>(
*file_block->get_by_position(result_column_id).column)
.get_data();
@@ -702,11 +703,12 @@ Status execute_delete_conjuncts(const
format::FileScanRequest& request, int64_t
break;
}
DORIS_CHECK(delete_conjunct != nullptr);
+ const size_t original_columns = file_block->columns();
int result_column_id = -1;
RETURN_IF_ERROR(delete_conjunct->root()->execute(delete_conjunct.get(),
file_block,
&result_column_id));
- DORIS_CHECK(result_column_id >= 0 &&
- result_column_id <
static_cast<int>(file_block->columns()));
+ RETURN_IF_ERROR(detail::validate_ephemeral_expr_result_column(
+ original_columns, result_column_id, file_block->columns()));
const auto& delete_filter = assert_cast<const ColumnUInt8&>(
*file_block->get_by_position(result_column_id).column)
.get_data();
@@ -727,6 +729,19 @@ Status execute_delete_conjuncts(const
format::FileScanRequest& request, int64_t
} // namespace
+Status detail::validate_ephemeral_expr_result_column(size_t original_columns,
int result_column_id,
+ size_t current_columns) {
+ // Delete predicates may erase only a temporary expression result. A bare
SlotRef returns an
+ // input column id, which must remain in the block for later predicates
and materialization.
+ if (UNLIKELY(result_column_id < 0 || static_cast<size_t>(result_column_id)
< original_columns ||
+ static_cast<size_t>(result_column_id) >= current_columns)) {
+ return Status::InternalError(
+ "Delete conjunct result column {} is not ephemeral
(original={}, current={})",
+ result_column_id, original_columns, current_columns);
+ }
+ return Status::OK();
+}
+
uint16_t apply_compact_filter_to_selection(const IColumn::Filter& filter,
SelectionVector* selection,
uint16_t selected_rows) {
DORIS_CHECK(selection != nullptr);
diff --git a/be/src/format_v2/parquet/parquet_scan.h
b/be/src/format_v2/parquet/parquet_scan.h
index dc926c3a766..ac883cbbcbe 100644
--- a/be/src/format_v2/parquet/parquet_scan.h
+++ b/be/src/format_v2/parquet/parquet_scan.h
@@ -77,6 +77,8 @@ std::vector<size_t> adaptive_prefetch_prefix(
const std::unordered_map<size_t, AdaptivePredicateStats>& stats,
double minimum_reach_probability);
bool should_sample_adaptive_predicate(size_t samples, size_t batch_sequence);
+Status validate_ephemeral_expr_result_column(size_t original_columns, int
result_column_id,
+ size_t current_columns);
Status build_native_prefetch_ranges(
const tparquet::FileMetaData& metadata,
const std::vector<std::unique_ptr<ParquetColumnSchema>>& file_schema,
diff --git a/be/src/format_v2/parquet/reader/native/column_chunk_reader.cpp
b/be/src/format_v2/parquet/reader/native/column_chunk_reader.cpp
index c219a7aef03..26afaf15bf4 100644
--- a/be/src/format_v2/parquet/reader/native/column_chunk_reader.cpp
+++ b/be/src/format_v2/parquet/reader/native/column_chunk_reader.cpp
@@ -21,6 +21,7 @@
#include <gen_cpp/parquet_types.h>
#include <glog/logging.h>
#include <parquet/metadata.h>
+#include <snappy.h>
#include <string.h>
#include <algorithm>
@@ -86,6 +87,102 @@ Status validate_uncompressed_page_sizes(const
tparquet::PageHeader& header,
return Status::OK();
}
+Status validate_fixed_width_page_size(const tparquet::PageHeader& header,
int32_t type_length,
+ level_t max_rep_level, level_t
max_def_level,
+ bool schema_is_required) {
+ if (type_length <= 0) {
+ return Status::OK();
+ }
+ const bool is_v2 = header.__isset.data_page_header_v2;
+ if (!is_v2 && !header.__isset.data_page_header) {
+ return Status::OK();
+ }
+ const auto encoding =
+ is_v2 ? header.data_page_header_v2.encoding :
header.data_page_header.encoding;
+ if (encoding != tparquet::Encoding::PLAIN &&
+ encoding != tparquet::Encoding::BYTE_STREAM_SPLIT) {
+ return Status::OK();
+ }
+ int32_t num_physical_values = 0;
+ int64_t level_bytes = 0;
+ if (is_v2) {
+ const auto& page = header.data_page_header_v2;
+ if (UNLIKELY(page.num_values < 0 || page.num_nulls < 0 ||
+ page.num_nulls > page.num_values ||
page.repetition_levels_byte_length < 0 ||
+ page.definition_levels_byte_length < 0)) {
+ return Status::Corruption("Parquet data page v2 has invalid value
or level counts");
+ }
+ num_physical_values = page.num_values - page.num_nulls;
+ level_bytes = static_cast<int64_t>(page.repetition_levels_byte_length)
+
+ page.definition_levels_byte_length;
+ } else {
+ if (max_rep_level != 0 || max_def_level != 0 || !schema_is_required) {
+ return Status::OK();
+ }
+ num_physical_values = header.data_page_header.num_values;
+ }
+ if (level_bytes > std::numeric_limits<int32_t>::max() ||
num_physical_values < 0 ||
+ static_cast<uint64_t>(num_physical_values) >
+ (static_cast<uint64_t>(std::numeric_limits<int32_t>::max()) -
level_bytes) /
+ static_cast<uint32_t>(type_length)) {
+ return Status::Corruption("Parquet fixed-width PLAIN page byte size
overflows");
+ }
+ const int64_t expected = level_bytes +
static_cast<int64_t>(num_physical_values) * type_length;
+ if (UNLIKELY(header.uncompressed_page_size != expected)) {
+ // V2 exposes null and level extents separately, so fixed-width
payload size is known before
+ // decompression even for optional columns and must gate
attacker-controlled allocation.
+ return Status::Corruption("Parquet fixed-width page has {}
uncompressed bytes, expected {}",
+ header.uncompressed_page_size, expected);
+ }
+ return Status::OK();
+}
+
+Status validate_dictionary_page_size(const tparquet::PageHeader& header,
int32_t type_length) {
+ DORIS_CHECK(header.__isset.dictionary_page_header);
+ const int32_t num_values = header.dictionary_page_header.num_values;
+ if (UNLIKELY(num_values < 0 || (num_values == 0 &&
header.uncompressed_page_size != 0))) {
+ // An empty dictionary owns no payload; validate before allocating
from its untrusted size.
+ return Status::Corruption("Parquet dictionary has {} values and {}
uncompressed bytes",
+ num_values, header.uncompressed_page_size);
+ }
+ if (type_length > 0) {
+ if (UNLIKELY(static_cast<uint64_t>(num_values) >
+
static_cast<uint64_t>(std::numeric_limits<int32_t>::max()) /
+ static_cast<uint32_t>(type_length))) {
+ return Status::Corruption("Parquet fixed-width dictionary byte
size overflows");
+ }
+ const int64_t expected = static_cast<int64_t>(num_values) *
type_length;
+ if (UNLIKELY(header.uncompressed_page_size != expected)) {
+ // Fixed-width dictionaries have no level section, so reject
forged extents before the
+ // decoder allocates storage based on the untrusted page header.
+ return Status::Corruption(
+ "Parquet fixed-width dictionary has {} uncompressed bytes,
expected {}",
+ header.uncompressed_page_size, expected);
+ }
+ }
+ return Status::OK();
+}
+
+Status validate_compressed_page_size(tparquet::CompressionCodec::type codec,
+ const Slice& compressed_data,
+ size_t expected_uncompressed_size) {
+ if (codec != tparquet::CompressionCodec::SNAPPY) {
+ return Status::OK();
+ }
+ size_t actual_uncompressed_size = 0;
+ if (UNLIKELY(!snappy::GetUncompressedLength(compressed_data.data,
compressed_data.size,
+ &actual_uncompressed_size))) {
+ return Status::Corruption("Invalid Snappy-compressed Parquet page");
+ }
+ if (UNLIKELY(actual_uncompressed_size != expected_uncompressed_size)) {
+ // Snappy exposes its decoded extent without an output buffer. Check
it before trusting the
+ // page header so malformed variable-width pages cannot force a
header-sized allocation.
+ return Status::Corruption("Snappy Parquet page expands to {} bytes,
expected {}",
+ actual_uncompressed_size,
expected_uncompressed_size);
+ }
+ return Status::OK();
+}
+
ParquetReaderCompat parquet_reader_compat(const std::string& created_by) {
if (created_by.empty()) {
return {};
@@ -905,6 +1002,18 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::parse_page_header() {
template <bool IN_COLLECTION, bool OFFSET_INDEX>
Status ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>::next_page() {
+ // Level parsing advances _page_data past the allocation base, so retain
explicit ownership
+ // state instead of inferring whether current decoders still reference
decompressed storage.
+ _page_uses_decompress_buf = false;
+ _active_decompress_bytes = 0;
+ if (_decompress_release_pending) {
+ if (_decompress_buf_size > _decompress_release_threshold) {
+ _decompress_buf.reset();
+ _decompress_buf_size = 0;
+ }
+ _decompress_release_pending = false;
+ _decompress_release_threshold = std::numeric_limits<size_t>::max();
+ }
_state = INITIALIZED;
RETURN_IF_ERROR(_page_reader->next_page());
return Status::OK();
@@ -939,8 +1048,17 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::load_page_data() {
RETURN_IF_ERROR(_page_reader->get_page_header(&header));
RETURN_IF_ERROR(validate_uncompressed_page_sizes(
*header, _metadata.codec,
_page_read_ctx.data_page_v2_always_compressed));
+ // Zero levels alone are insufficient: test/protocol adapters can leave
repetition unset, so
+ // only an explicitly REQUIRED schema proves that every logical value has
fixed-width bytes.
+ const bool schema_is_required =
_field_schema->parquet_schema.__isset.repetition_type &&
+
_field_schema->parquet_schema.repetition_type ==
+
tparquet::FieldRepetitionType::REQUIRED;
+ RETURN_IF_ERROR(validate_fixed_width_page_size(*header,
_get_type_length(), _max_rep_level,
+ _max_def_level,
schema_is_required));
int32_t uncompressed_size = header->uncompressed_page_size;
bool page_loaded = false;
+ _page_uses_decompress_buf = false;
+ _active_decompress_bytes = 0;
// First, try to reuse a cache handle previously discovered by PageReader
// (header-only lookup) to avoid a second lookup here.
@@ -992,8 +1110,12 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::load_page_data() {
header->__isset.data_page_header_v2
?
static_cast<size_t>(header->uncompressed_page_size) - levels_size
:
static_cast<size_t>(header->uncompressed_page_size);
+ RETURN_IF_ERROR(validate_compressed_page_size(_metadata.codec,
payload_slice,
+
uncompressed_payload_size));
_reserve_decompress_buf(uncompressed_payload_size);
_page_data = Slice(_decompress_buf.get(),
uncompressed_payload_size);
+ _page_uses_decompress_buf = true;
+ _active_decompress_bytes = uncompressed_payload_size;
SCOPED_RAW_TIMER(&_chunk_statistics.decompress_time);
_chunk_statistics.decompress_cnt++;
RETURN_IF_ERROR(_block_compress_codec->decompress(payload_slice, &_page_data));
@@ -1040,8 +1162,12 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::load_page_data() {
if (page_has_compression) {
// Decompress payload for immediate decoding
+ RETURN_IF_ERROR(validate_compressed_page_size(
+ _metadata.codec, compressed_data,
static_cast<size_t>(uncompressed_size)));
_reserve_decompress_buf(uncompressed_size);
_page_data = Slice(_decompress_buf.get(), uncompressed_size);
+ _page_uses_decompress_buf = true;
+ _active_decompress_bytes =
static_cast<size_t>(uncompressed_size);
SCOPED_RAW_TIMER(&_chunk_statistics.decompress_time);
_chunk_statistics.decompress_cnt++;
RETURN_IF_ERROR(_block_compress_codec->decompress(compressed_data,
&_page_data));
@@ -1186,7 +1312,15 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::_decode_dict_page() {
int32_t uncompressed_size = header->uncompressed_page_size;
RETURN_IF_ERROR(validate_uncompressed_page_sizes(
*header, _metadata.codec,
_page_read_ctx.data_page_v2_always_compressed));
- auto dict_data = make_unique_buffer<uint8_t>(uncompressed_size);
+ RETURN_IF_ERROR(validate_dictionary_page_size(*header,
_get_type_length()));
+ DorisUniqueBufferPtr<uint8_t> dict_data;
+ bool dict_buffer_allocated = false;
+ const auto allocate_dict_buffer = [&]() {
+ if (!dict_buffer_allocated) {
+ dict_data = make_unique_buffer<uint8_t>(uncompressed_size);
+ dict_buffer_allocated = true;
+ }
+ };
bool dict_loaded = false;
// Try to load dictionary page from cache
@@ -1196,6 +1330,10 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::_decode_dict_page() {
const PageCacheHandle& handle = _page_reader->page_cache_handle();
Slice cached = handle.data();
size_t header_size = _page_reader->header_bytes().size();
+ if (UNLIKELY(header_size > cached.size)) {
+ return Status::Corruption(
+ "Cached Parquet dictionary is shorter than its page
header");
+ }
// Dictionary page layout in cache: header | payload (compressed
or uncompressed)
Slice payload_slice(cached.data + header_size, cached.size -
header_size);
@@ -1208,11 +1346,15 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::_decode_dict_page() {
"Cached Parquet dictionary payload has size {},
expected {}",
payload_slice.size, uncompressed_size);
}
+ allocate_dict_buffer();
memcpy(dict_data.get(), payload_slice.data,
payload_slice.size);
dict_loaded = true;
} else {
CHECK(_block_compress_codec);
// Decompress cached compressed dictionary data
+ RETURN_IF_ERROR(validate_compressed_page_size(
+ _metadata.codec, payload_slice,
static_cast<size_t>(uncompressed_size)));
+ allocate_dict_buffer();
Slice dict_slice(dict_data.get(), uncompressed_size);
{
SCOPED_RAW_TIMER(&_chunk_statistics.decompress_time);
@@ -1239,15 +1381,13 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::_decode_dict_page() {
// Load and decompress dictionary page from file
if (_block_compress_codec != nullptr) {
auto dict_num = header->dictionary_page_header.num_values;
- if (dict_num == 0 && uncompressed_size != 0) {
- return Status::IOError(
- "Dictionary page's num_values is {} but
uncompressed_size is {}", dict_num,
- uncompressed_size);
- }
Slice compressed_data;
- Slice dict_slice(dict_data.get(), uncompressed_size);
if (dict_num != 0) {
RETURN_IF_ERROR(_page_reader->get_page_data(compressed_data));
+ RETURN_IF_ERROR(validate_compressed_page_size(
+ _metadata.codec, compressed_data,
static_cast<size_t>(uncompressed_size)));
+ allocate_dict_buffer();
+ Slice dict_slice(dict_data.get(), uncompressed_size);
// Dictionary probes stop before data pages, so count their
decompression here or
// metadata pruning profiles will report no codec work for the
scan.
{
@@ -1262,6 +1402,8 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::_decode_dict_page() {
dict_slice.size, uncompressed_size);
}
}
+ allocate_dict_buffer();
+ Slice dict_slice(dict_data.get(), uncompressed_size);
// Decide whether to cache decompressed or compressed dictionary
based on threshold
// If uncompressed_page_size == 0, should_cache_decompressed will
return true
@@ -1295,6 +1437,11 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::_decode_dict_page() {
} else {
Slice dict_slice;
RETURN_IF_ERROR(_page_reader->get_page_data(dict_slice));
+ if (UNLIKELY(dict_slice.size !=
static_cast<size_t>(uncompressed_size))) {
+ return Status::Corruption("Parquet dictionary payload has size
{}, expected {}",
+ dict_slice.size, uncompressed_size);
+ }
+ allocate_dict_buffer();
// The data is stored by BufferedStreamReader, we should copy it
out
memcpy(dict_data.get(), dict_slice.data, dict_slice.size);
@@ -1309,6 +1456,7 @@ Status ColumnChunkReader<IN_COLLECTION,
OFFSET_INDEX>::_decode_dict_page() {
}
}
}
+ allocate_dict_buffer();
// Cache page decoder
std::unique_ptr<Decoder> page_decoder;
diff --git a/be/src/format_v2/parquet/reader/native/column_chunk_reader.h
b/be/src/format_v2/parquet/reader/native/column_chunk_reader.h
index 5acf544bcc9..102366f2494 100644
--- a/be/src/format_v2/parquet/reader/native/column_chunk_reader.h
+++ b/be/src/format_v2/parquet/reader/native/column_chunk_reader.h
@@ -19,8 +19,10 @@
#include <gen_cpp/parquet_types.h>
+#include <algorithm>
#include <cstddef>
#include <cstdint>
+#include <limits>
#include <memory>
#include <string>
#include <unordered_map>
@@ -71,6 +73,13 @@ bool can_prepare_page_cache_payload(bool
session_cache_enabled, bool storage_cac
Status validate_uncompressed_page_sizes(const tparquet::PageHeader& header,
tparquet::CompressionCodec::type codec,
bool data_page_v2_always_compressed);
+Status validate_fixed_width_page_size(const tparquet::PageHeader& header,
int32_t type_length,
+ level_t max_rep_level, level_t
max_def_level,
+ bool schema_is_required = true);
+Status validate_dictionary_page_size(const tparquet::PageHeader& header,
int32_t type_length = -1);
+Status validate_compressed_page_size(tparquet::CompressionCodec::type codec,
+ const Slice& compressed_data,
+ size_t expected_uncompressed_size);
struct ColumnChunkReaderStatistics {
int64_t decompress_time = 0;
@@ -215,10 +224,24 @@ public:
// Level decoders may batch-convert unsigned RLE values into Doris'
signed level_t.
_rep_level_decoder.release_scratch(max_retained_bytes);
_def_level_decoder.release_scratch(max_retained_bytes);
+ if (_decompress_buf_size > max_retained_bytes) {
+ if (_page_uses_decompress_buf) {
+ // Keep the request until the page boundary because decoders
still point into this
+ // allocation; dropping it now would trade retained memory for
a use-after-free.
+ _decompress_release_pending = true;
+ _decompress_release_threshold =
+ std::min(_decompress_release_threshold,
max_retained_bytes);
+ } else {
+ _decompress_buf.reset();
+ _decompress_buf_size = 0;
+ _decompress_release_pending = false;
+ _decompress_release_threshold =
std::numeric_limits<size_t>::max();
+ }
+ }
}
size_t retained_decoder_scratch_bytes() const {
- size_t bytes = _rep_level_decoder.retained_scratch_bytes() +
+ size_t bytes = _decompress_buf_size +
_rep_level_decoder.retained_scratch_bytes() +
_def_level_decoder.retained_scratch_bytes();
for (const auto& [encoding, decoder] : _decoders) {
bytes += decoder->retained_scratch_bytes();
@@ -229,7 +252,8 @@ public:
size_t active_decoder_scratch_bytes() const {
// Only the current encoding is active. Old decoder instances retain
reusable capacity but
// must not make the high-water policy treat their last batch as
current working memory.
- return (_page_decoder == nullptr ? 0 :
_page_decoder->active_scratch_bytes()) +
+ return _active_decompress_bytes +
+ (_page_decoder == nullptr ? 0 :
_page_decoder->active_scratch_bytes()) +
_rep_level_decoder.active_scratch_bytes() +
_def_level_decoder.active_scratch_bytes();
}
@@ -367,6 +391,10 @@ private:
Slice _page_data;
DorisUniqueBufferPtr<uint8_t> _decompress_buf;
size_t _decompress_buf_size = 0;
+ bool _page_uses_decompress_buf = false;
+ size_t _active_decompress_bytes = 0;
+ bool _decompress_release_pending = false;
+ size_t _decompress_release_threshold = std::numeric_limits<size_t>::max();
Slice _v2_rep_levels;
Slice _v2_def_levels;
bool _dict_checked = false;
diff --git a/be/src/format_v2/parquet/reader/native/column_reader.cpp
b/be/src/format_v2/parquet/reader/native/column_reader.cpp
index 9417da15683..592e048846e 100644
--- a/be/src/format_v2/parquet/reader/native/column_reader.cpp
+++ b/be/src/format_v2/parquet/reader/native/column_reader.cpp
@@ -1271,6 +1271,10 @@ ScalarColumnReader<IN_COLLECTION,
OFFSET_INDEX>::materialize_dictionary_values(
// interpretation as ordinary data-page decoding.
const DataTypePtr dictionary_type = remove_nullable(target_type);
const DataTypeSerDeSPtr dictionary_serde = dictionary_type->get_serde();
+ // The probe is the first typed read for this reader. Publish its SerDe
identity now so the
+ // first dictionary-ID batch does not mistake initialization for a type
change and drop cache.
+ _serde_type = dictionary_type.get();
+ _serde = dictionary_serde;
if (_materialization_state.dictionary_generation !=
dictionary_decoder->dictionary_generation()) {
_materialization_state.typed_dictionary =
dictionary_type->create_column();
@@ -1285,6 +1289,9 @@ ScalarColumnReader<IN_COLLECTION,
OFFSET_INDEX>::materialize_dictionary_values(
DORIS_CHECK_EQ(_materialization_state.typed_dictionary->size(),
dictionary_decoder->dictionary_size());
_materialization_state.dictionary_generation =
dictionary_decoder->dictionary_generation();
+#ifdef BE_TEST
+ ++_dictionary_materialization_count;
+#endif
}
auto result = _materialization_state.typed_dictionary->clone_empty();
@@ -1350,11 +1357,16 @@ Status ScalarColumnReader<IN_COLLECTION,
OFFSET_INDEX>::read_column_data(
}
const DataTypePtr serde_type = is_dict_filter ? file_type :
materialization_type;
- if (_serde_type != serde_type.get() || _dictionary_index_only !=
is_dict_filter) {
+ const bool serde_type_changed = _serde_type != serde_type.get();
+ if (serde_type_changed || _dictionary_index_only != is_dict_filter) {
_serde_type = serde_type.get();
_serde = serde_type->get_serde();
_dictionary_index_only = is_dict_filter;
- _materialization_state.reset_dictionary();
+ // Switching between typed values and dictionary IDs does not
invalidate the dictionary
+ // materialized by the pruning probe. Preserve it unless the logical
SerDe type changed.
+ if (serde_type_changed) {
+ _materialization_state.reset_dictionary();
+ }
}
_decode_context.dictionary_index_only = is_dict_filter;
diff --git a/be/src/format_v2/parquet/reader/native/column_reader.h
b/be/src/format_v2/parquet/reader/native/column_reader.h
index b9a25a8132e..6ea2787f9f0 100644
--- a/be/src/format_v2/parquet/reader/native/column_reader.h
+++ b/be/src/format_v2/parquet/reader/native/column_reader.h
@@ -300,6 +300,9 @@ public:
#ifdef BE_TEST
void reserve_batch_scratch_for_test(size_t elements);
size_t retained_batch_scratch_bytes_for_test() const;
+ size_t dictionary_materialization_count_for_test() const {
+ return _dictionary_materialization_count;
+ }
#endif
void reset_filter_map_index() override {
@@ -311,6 +314,9 @@ private:
const tparquet::OffsetIndex* _offset_index = nullptr;
std::unique_ptr<io::BufferedFileStreamReader> _stream_reader;
std::unique_ptr<ColumnChunkReader<IN_COLLECTION, OFFSET_INDEX>>
_chunk_reader;
+#ifdef BE_TEST
+ size_t _dictionary_materialization_count = 0;
+#endif
// rep def levels buffer.
std::vector<level_t> _rep_levels;
std::vector<level_t> _def_levels;
diff --git a/be/src/format_v2/parquet/reader/native/delta_bit_pack_decoder.h
b/be/src/format_v2/parquet/reader/native/delta_bit_pack_decoder.h
index 531cb770407..db9efecaced 100644
--- a/be/src/format_v2/parquet/reader/native/delta_bit_pack_decoder.h
+++ b/be/src/format_v2/parquet/reader/native/delta_bit_pack_decoder.h
@@ -672,6 +672,12 @@ Status DeltaBitPackDecoder<T>::_init_header() {
if (_mini_blocks_per_block == 0) {
return Status::InvalidArgument("Cannot have zero miniblock per block");
}
+ // Parquet requires integral block geometry; truncating here would
silently decode fewer
+ // values than the header declares and desynchronize the following block.
+ if (UNLIKELY(_values_per_block % _mini_blocks_per_block != 0)) {
+ return Status::Corruption("Parquet delta block size {} is not
divisible by {} miniblocks",
+ _values_per_block, _mini_blocks_per_block);
+ }
_values_per_mini_block = _values_per_block / _mini_blocks_per_block;
if (_values_per_mini_block == 0) {
return Status::InvalidArgument("Cannot have zero value per miniblock");
diff --git a/be/src/format_v2/parquet/reader/native_column_reader.cpp
b/be/src/format_v2/parquet/reader/native_column_reader.cpp
index 8aa79ecb7e7..233b37c096e 100644
--- a/be/src/format_v2/parquet/reader/native_column_reader.cpp
+++ b/be/src/format_v2/parquet/reader/native_column_reader.cpp
@@ -323,13 +323,7 @@ Status NativeColumnReader::read_with_filter(int64_t rows,
const uint8_t* filter_
if (_nested && _profile.nested_batches != nullptr) {
COUNTER_UPDATE(_profile.nested_batches, 1);
}
- // Retained-capacity inspection walks the native reader tree. Check it
periodically instead of
- // on every small batch; row-group destruction is still the hard lifetime
bound for scratch.
- constexpr size_t SCRATCH_CHECK_BATCH_INTERVAL = 16;
- if (++_batches_since_scratch_check >= SCRATCH_CHECK_BATCH_INTERVAL) {
-
_native_reader->release_batch_scratch(MAX_RETAINED_BATCH_SCRATCH_BYTES);
- _batches_since_scratch_check = 0;
- }
+ release_batch_scratch_if_needed();
if (*rows_read != rows) {
return Status::Corruption("Native parquet reader returned {} rows,
expected {} for {}",
*rows_read, rows, _name);
@@ -397,9 +391,20 @@ Status
NativeColumnReader::read_with_fixed_width_filter(int64_t rows, const uint
*rows_read, rows, _name);
}
*used_filter = true;
+ release_batch_scratch_if_needed();
return Status::OK();
}
+void NativeColumnReader::release_batch_scratch_if_needed() {
+ // PLAIN predicate batches bypass materialization but share the same
persistent decoder tree,
+ // so both read paths must advance the retained-capacity aging clock.
+ constexpr size_t SCRATCH_CHECK_BATCH_INTERVAL = 16;
+ if (++_batches_since_scratch_check >= SCRATCH_CHECK_BATCH_INTERVAL) {
+
_native_reader->release_batch_scratch(MAX_RETAINED_BATCH_SCRATCH_BYTES);
+ _batches_since_scratch_check = 0;
+ }
+}
+
Status NativeColumnReader::validate_selected_span(int64_t rows) {
DORIS_CHECK(rows >= 0);
while (_selected_range_idx < _selected_ranges.size()) {
diff --git a/be/src/format_v2/parquet/reader/native_column_reader.h
b/be/src/format_v2/parquet/reader/native_column_reader.h
index 0ffb8d05400..1d92735bfa1 100644
--- a/be/src/format_v2/parquet/reader/native_column_reader.h
+++ b/be/src/format_v2/parquet/reader/native_column_reader.h
@@ -116,6 +116,7 @@ private:
const VExprSPtrs& conjuncts, int
column_id,
IColumn* projected_column,
IColumn::Filter* row_filter,
int64_t* rows_read, bool* used_filter);
+ void release_batch_scratch_if_needed();
int64_t sync_native_profile();
void record_page_fragments(int64_t page_fragments);
Status validate_selected_span(int64_t rows);
diff --git a/be/src/format_v2/table/hudi_reader.cpp
b/be/src/format_v2/table/hudi_reader.cpp
index ef9c2772543..271f14a7b48 100644
--- a/be/src/format_v2/table/hudi_reader.cpp
+++ b/be/src/format_v2/table/hudi_reader.cpp
@@ -146,6 +146,13 @@ const format::MaterializedBlockStats&
HudiHybridReader::last_materialized_block_
:
format::TableReader::last_materialized_block_stats();
}
+int64_t HudiHybridReader::condition_cache_hit_count() const {
+ // Keep the wrapper count cumulative across native/JNI dispatch so
scanner-level delta
+ // accounting neither loses a child hit nor observes a counter reset on a
split switch.
+ return (_native_reader == nullptr ? 0 :
_native_reader->condition_cache_hit_count()) +
+ (_jni_reader == nullptr ? 0 :
_jni_reader->condition_cache_hit_count());
+}
+
Status HudiHybridReader::_ensure_current_split_reader(const
format::SplitReadOptions& options) {
DORIS_CHECK(_scan_params != nullptr);
if (_is_jni_split(*_scan_params, options.current_range)) {
diff --git a/be/src/format_v2/table/hudi_reader.h
b/be/src/format_v2/table/hudi_reader.h
index 42776422c4a..c13ac1215d7 100644
--- a/be/src/format_v2/table/hudi_reader.h
+++ b/be/src/format_v2/table/hudi_reader.h
@@ -67,6 +67,7 @@ public:
void set_batch_size(size_t batch_size) override;
Status append_conjuncts(const VExprContextSPtrs& conjuncts) override;
const format::MaterializedBlockStats& last_materialized_block_stats()
const override;
+ int64_t condition_cache_hit_count() const override;
#ifdef BE_TEST
void TEST_install_batch_size_children() {
@@ -76,6 +77,10 @@ public:
std::pair<size_t, size_t> TEST_child_batch_sizes() const {
return {_native_reader->TEST_batch_size(),
_jni_reader->TEST_batch_size()};
}
+ void TEST_set_child_condition_cache_hits(int64_t native_hits, int64_t
jni_hits) {
+ _native_reader->TEST_set_condition_cache_hit_count(native_hits);
+ _jni_reader->TEST_set_condition_cache_hit_count(jni_hits);
+ }
void TEST_set_child_reader_factories(
std::function<std::unique_ptr<format::TableReader>()>
native_factory,
std::function<std::unique_ptr<format::TableReader>()> jni_factory)
{
diff --git
a/be/src/format_v2/table/iceberg_position_delete_sys_table_reader.cpp
b/be/src/format_v2/table/iceberg_position_delete_sys_table_reader.cpp
index 16c7f829089..6ab67caa77d 100644
--- a/be/src/format_v2/table/iceberg_position_delete_sys_table_reader.cpp
+++ b/be/src/format_v2/table/iceberg_position_delete_sys_table_reader.cpp
@@ -146,6 +146,7 @@ protected:
}
void configure_mapper_options(format::TableColumnMapperOptions* options)
const override {
+ options->enable_row_lineage_virtual_columns = true;
// Parquet may preserve a selected complex wrapper without its own ID;
position-delete row
// projection must use the same descendant-ID fallback as ordinary
Iceberg data scans.
options->allow_idless_complex_wrapper_projection =
diff --git a/be/src/format_v2/table/iceberg_reader.h
b/be/src/format_v2/table/iceberg_reader.h
index cf864188074..d28be3d7f98 100644
--- a/be/src/format_v2/table/iceberg_reader.h
+++ b/be/src/format_v2/table/iceberg_reader.h
@@ -69,6 +69,7 @@ public:
protected:
void configure_mapper_options(format::TableColumnMapperOptions* options)
const override {
+ options->enable_row_lineage_virtual_columns = true;
options->allow_idless_complex_wrapper_projection =
supports_iceberg_scan_semantics_v1(_scan_params) && _format ==
FileFormat::PARQUET;
}
diff --git a/be/src/format_v2/table/paimon_reader.cpp
b/be/src/format_v2/table/paimon_reader.cpp
index 3e409fee9da..d815ef81c9b 100644
--- a/be/src/format_v2/table/paimon_reader.cpp
+++ b/be/src/format_v2/table/paimon_reader.cpp
@@ -173,6 +173,13 @@ const format::MaterializedBlockStats&
PaimonHybridReader::last_materialized_bloc
:
format::TableReader::last_materialized_block_stats();
}
+int64_t PaimonHybridReader::condition_cache_hit_count() const {
+ // Both children survive split switches, so the wrapper must publish their
cumulative totals;
+ // returning only the active child would make FileScannerV2's monotonic
delta go backwards.
+ return (_native_reader == nullptr ? 0 :
_native_reader->condition_cache_hit_count()) +
+ (_jni_reader == nullptr ? 0 :
_jni_reader->condition_cache_hit_count());
+}
+
Status PaimonHybridReader::_ensure_current_split_reader(const
format::SplitReadOptions& options) {
if (_is_jni_split(options.current_range)) {
DCHECK(options.current_split_format == format::FileFormat::JNI);
diff --git a/be/src/format_v2/table/paimon_reader.h
b/be/src/format_v2/table/paimon_reader.h
index 6ee5e7d72f3..b4f076fc6e4 100644
--- a/be/src/format_v2/table/paimon_reader.h
+++ b/be/src/format_v2/table/paimon_reader.h
@@ -73,6 +73,7 @@ public:
void set_batch_size(size_t batch_size) override;
Status append_conjuncts(const VExprContextSPtrs& conjuncts) override;
const format::MaterializedBlockStats& last_materialized_block_stats()
const override;
+ int64_t condition_cache_hit_count() const override;
#ifdef BE_TEST
static bool TEST_is_jni_split(const TFileRangeDesc& range) { return
_is_jni_split(range); }
@@ -87,6 +88,10 @@ public:
std::pair<size_t, size_t> TEST_child_batch_sizes() const {
return {_native_reader->TEST_batch_size(),
_jni_reader->TEST_batch_size()};
}
+ void TEST_set_child_condition_cache_hits(int64_t native_hits, int64_t
jni_hits) {
+ _native_reader->TEST_set_condition_cache_hit_count(native_hits);
+ _jni_reader->TEST_set_condition_cache_hit_count(jni_hits);
+ }
void TEST_set_child_reader_factories(
std::function<std::unique_ptr<format::TableReader>()>
native_factory,
std::function<std::unique_ptr<format::TableReader>()> jni_factory)
{
diff --git a/be/src/format_v2/table/remote_doris_reader.cpp
b/be/src/format_v2/table/remote_doris_reader.cpp
index db278d9b0ce..bd2c08b5f23 100644
--- a/be/src/format_v2/table/remote_doris_reader.cpp
+++ b/be/src/format_v2/table/remote_doris_reader.cpp
@@ -20,8 +20,13 @@
#include <arrow/flight/client.h>
#include <arrow/flight/types.h>
+#include <algorithm>
+#include <chrono>
+#include <condition_variable>
#include <memory>
+#include <mutex>
#include <string>
+#include <thread>
#include <utility>
#include <vector>
@@ -38,7 +43,9 @@
#include "format_v2/materialized_reader_util.h"
#include "runtime/descriptors.h"
#include "runtime/file_scan_profile.h"
+#include "runtime/query_context.h"
#include "runtime/runtime_state.h"
+#include "runtime/thread_context.h"
#include "util/timezone_utils.h"
namespace doris::format::remote_doris {
@@ -64,7 +71,12 @@ Status validate_remote_doris_range(const TFileRangeDesc&
range) {
class FlightRemoteDorisStream final : public RemoteDorisStream {
public:
- explicit FlightRemoteDorisStream(const TFileRangeDesc& range) :
_range(range) {}
+ FlightRemoteDorisStream(const TFileRangeDesc& range,
std::shared_ptr<io::IOContext> io_ctx,
+ RuntimeState* runtime_state, int timeout_seconds)
+ : _range(range),
+ _io_ctx(std::move(io_ctx)),
+ _runtime_state(runtime_state),
+ _timeout_seconds(std::max(1, timeout_seconds)) {}
Status open() {
RETURN_IF_ERROR(validate_remote_doris_range(_range));
@@ -75,14 +87,103 @@ public:
arrow::flight::Ticket ticket;
RETURN_DORIS_STATUS_IF_ERROR(
arrow::flight::Ticket::Deserialize(params.ticket).Value(&ticket));
+ struct PendingOpen {
+ std::mutex mutex;
+ std::condition_variable cv;
+ bool done = false;
+ bool abandoned = false;
+ arrow::Status status = arrow::Status::OK();
+ std::unique_ptr<arrow::flight::FlightClient> client;
+ std::unique_ptr<arrow::flight::FlightStreamReader> stream;
+ };
+ auto pending = std::make_shared<PendingOpen>();
+ std::unique_ptr<arrow::flight::FlightClient> flight_client;
RETURN_DORIS_STATUS_IF_ERROR(
-
arrow::flight::FlightClient::Connect(location).Value(&_flight_client));
-
RETURN_DORIS_STATUS_IF_ERROR(_flight_client->DoGet(ticket).Value(&_stream));
+
arrow::flight::FlightClient::Connect(location).Value(&flight_client));
+ arrow::flight::FlightCallOptions options;
+ // A Flight deadline covers streaming reads as well as DoGet setup, so
a stalled Next()
+ // cannot outlive the query execution timeout indefinitely.
+ options.timeout = std::chrono::seconds(_timeout_seconds);
+ // Start before DoGet because endpoint setup is itself a blocking RPC
covered by the same
+ // query/scanner cancellation contract as streaming Next().
+ _cancellation_watcher = std::jthread(
+ [this](std::stop_token stop_token) {
_watch_cancellation(stop_token); });
+
+ std::shared_ptr<ResourceContext> resource_ctx;
+ if (_runtime_state != nullptr && _runtime_state->get_query_ctx() !=
nullptr) {
+ resource_ctx = _runtime_state->get_query_ctx()->resource_ctx();
+ }
+ std::thread do_get_thread([pending, options, ticket, resource_ctx,
+ client = std::move(flight_client)]()
mutable {
+ const auto do_get = [&] {
+ std::unique_ptr<arrow::flight::FlightStreamReader> stream;
+ auto status = client->DoGet(options, ticket).Value(&stream);
+ {
+ std::lock_guard lock(pending->mutex);
+ if (!pending->abandoned) {
+ pending->status = std::move(status);
+ pending->client = std::move(client);
+ pending->stream = std::move(stream);
+ } else {
+ // A detached worker must release its query-owned
Flight client
+ // before leaving the task attachment that accounts
for it.
+ client.reset();
+ }
+ pending->done = true;
+ }
+ pending->cv.notify_all();
+ };
+ if (resource_ctx != nullptr) {
+ SCOPED_ATTACH_TASK(resource_ctx);
+ do_get();
+ } else {
+ SCOPED_INIT_THREAD_CONTEXT();
+ do_get();
+ }
+ });
+ bool cancelled_during_open = false;
+ {
+ std::unique_lock lock(pending->mutex);
+ while (!pending->done && !_is_cancelled()) {
+ pending->cv.wait_for(lock, std::chrono::milliseconds(25));
+ }
+ if (!pending->done) {
+ pending->abandoned = true;
+ cancelled_during_open = true;
+ }
+ }
+ if (cancelled_during_open) {
+ // Arrow 17 exposes no cancellable handle until DoGet returns.
Detaching the bounded RPC
+ // keeps query/scanner shutdown prompt while the call is still
capped by its deadline;
+ // the shared state owns all Arrow objects until that worker exits.
+ do_get_thread.detach();
+ _stop_cancellation_watcher();
+ return Status::Cancelled("Remote Doris Flight open was cancelled");
+ }
+ do_get_thread.join();
+ if (!pending->status.ok()) {
+ _stop_cancellation_watcher();
+ RETURN_DORIS_STATUS_IF_ERROR(pending->status);
+ }
+ {
+ std::lock_guard lock(_flight_mutex);
+ _flight_client = std::move(pending->client);
+ _stream = std::move(pending->stream);
+ }
+ if (_is_cancelled()) {
+ _cancel_flight_call();
+ _stop_cancellation_watcher();
+ return Status::Cancelled("Remote Doris Flight open was cancelled");
+ }
return Status::OK();
}
Status next(std::shared_ptr<arrow::RecordBatch>* batch) override {
DORIS_CHECK(batch != nullptr);
+ if (_io_ctx != nullptr && _io_ctx->should_stop) {
+ _cancel_flight_call();
+ return Status::Cancelled("Remote Doris Flight read was cancelled");
+ }
arrow::flight::FlightStreamChunk chunk;
RETURN_DORIS_STATUS_IF_ERROR(_stream->Next().Value(&chunk));
*batch = chunk.data;
@@ -90,6 +191,13 @@ public:
}
Status close() override {
+ _stop_cancellation_watcher();
+ {
+ std::lock_guard lock(_flight_mutex);
+ if (_stream != nullptr) {
+ _stream->Cancel();
+ }
+ }
_stream.reset();
if (_flight_client != nullptr) {
RETURN_DORIS_STATUS_IF_ERROR(_flight_client->Close());
@@ -99,14 +207,65 @@ public:
}
private:
+ bool _is_cancelled() const {
+ return (_runtime_state != nullptr && _runtime_state->is_cancelled()) ||
+ (_io_ctx != nullptr && _io_ctx->should_stop);
+ }
+
+ void _cancel_flight_call() {
+ std::lock_guard lock(_flight_mutex);
+ if (_stream != nullptr) {
+ _stream->Cancel();
+ }
+ }
+
+ void _watch_cancellation_loop(std::stop_token watcher_stop_token) {
+ while (!watcher_stop_token.stop_requested()) {
+ if (_is_cancelled()) {
+ _cancel_flight_call();
+ return;
+ }
+ std::this_thread::sleep_for(std::chrono::milliseconds(25));
+ }
+ }
+
+ void _watch_cancellation(std::stop_token watcher_stop_token) {
+ if (_runtime_state != nullptr && _runtime_state->get_query_ctx() !=
nullptr &&
+ _runtime_state->get_query_ctx()->resource_ctx() != nullptr) {
+ // The watcher is query-owned and may allocate in Arrow while
signalling cancellation.
+ SCOPED_ATTACH_TASK(_runtime_state);
+ _watch_cancellation_loop(watcher_stop_token);
+ return;
+ }
+ // Metadata/tests can construct a RuntimeState without a QueryContext;
initialize TLS there
+ // instead of violating AttachTask's non-null resource-context
invariant.
+ SCOPED_INIT_THREAD_CONTEXT();
+ _watch_cancellation_loop(watcher_stop_token);
+ }
+
+ void _stop_cancellation_watcher() {
+ if (_cancellation_watcher.joinable()) {
+ _cancellation_watcher.request_stop();
+ _cancellation_watcher.join();
+ }
+ }
+
const TFileRangeDesc _range;
+ std::shared_ptr<io::IOContext> _io_ctx;
+ RuntimeState* _runtime_state;
+ int _timeout_seconds;
+ std::jthread _cancellation_watcher;
+ std::mutex _flight_mutex;
std::unique_ptr<arrow::flight::FlightClient> _flight_client;
std::unique_ptr<arrow::flight::FlightStreamReader> _stream;
};
-Status create_flight_stream(const TFileRangeDesc& range,
std::unique_ptr<RemoteDorisStream>* out) {
+Status create_flight_stream(const TFileRangeDesc& range,
std::shared_ptr<io::IOContext> io_ctx,
+ RuntimeState* runtime_state, int timeout_seconds,
+ std::unique_ptr<RemoteDorisStream>* out) {
DORIS_CHECK(out != nullptr);
- auto stream = std::make_unique<FlightRemoteDorisStream>(range);
+ auto stream = std::make_unique<FlightRemoteDorisStream>(range,
std::move(io_ctx), runtime_state,
+ timeout_seconds);
RETURN_IF_ERROR(stream->open());
*out = std::move(stream);
return Status::OK();
@@ -199,7 +358,10 @@ void RemoteDorisFileReader::_init_profile() {
Status RemoteDorisFileReader::init(RuntimeState* state) {
_init_profile();
SCOPED_TIMER(_total_time);
- (void)state;
+ if (state != nullptr) {
+ _flight_timeout_seconds = std::max(1, state->execution_timeout());
+ }
+ _runtime_state = state;
RETURN_IF_ERROR(validate_remote_doris_range(_range));
RETURN_IF_ERROR(_build_col_name_to_file_id());
_eof = false;
@@ -244,6 +406,14 @@ Status RemoteDorisFileReader::get_block(Block* file_block,
size_t* rows, bool* e
if (_stream == nullptr) {
return Status::InternalError("Remote Doris v2 reader is not open");
}
+ if (_io_ctx != nullptr && _io_ctx->should_stop) {
+ // Observe cancellation before entering a potentially blocking Flight
read; the production
+ // stream also carries a query-bounded RPC deadline for cancellation
arriving mid-read.
+ RETURN_IF_ERROR(close());
+ *rows = 0;
+ *eof = true;
+ return Status::OK();
+ }
*rows = 0;
*eof = false;
@@ -288,7 +458,8 @@ Status RemoteDorisFileReader::_open_stream() {
if (_stream_factory) {
RETURN_IF_ERROR(_stream_factory(_range, &_stream));
} else {
- RETURN_IF_ERROR(create_flight_stream(_range, &_stream));
+ RETURN_IF_ERROR(create_flight_stream(_range, _io_ctx, _runtime_state,
+ _flight_timeout_seconds,
&_stream));
}
DORIS_CHECK(_stream != nullptr);
return Status::OK();
diff --git a/be/src/format_v2/table/remote_doris_reader.h
b/be/src/format_v2/table/remote_doris_reader.h
index 79c37a29be4..3b4a03f4569 100644
--- a/be/src/format_v2/table/remote_doris_reader.h
+++ b/be/src/format_v2/table/remote_doris_reader.h
@@ -90,6 +90,8 @@ private:
RuntimeProfile::Counter* _io_time = nullptr;
RuntimeProfile::Counter* _materialize_time = nullptr;
RuntimeProfile::Counter* _filter_time = nullptr;
+ RuntimeState* _runtime_state = nullptr;
+ int _flight_timeout_seconds = 300;
std::unique_ptr<RemoteDorisStream> _stream;
std::unordered_map<std::string, LocalColumnId> _col_name_to_file_id;
};
diff --git a/be/src/format_v2/table_reader.h b/be/src/format_v2/table_reader.h
index 47d442f7fc1..48e3d996ce2 100644
--- a/be/src/format_v2/table_reader.h
+++ b/be/src/format_v2/table_reader.h
@@ -226,6 +226,7 @@ public:
size_t TEST_table_reader_owned_conjunct_count() const {
return _table_reader_owned_conjunct_count;
}
+ void TEST_set_condition_cache_hit_count(int64_t hits) {
_condition_cache_hit_count = hits; }
bool TEST_current_data_file_is_immutable() const {
DORIS_CHECK(_current_task != nullptr);
DORIS_CHECK(_current_task->data_file != nullptr);
@@ -404,7 +405,7 @@ public:
return Status::OK();
}
- int64_t condition_cache_hit_count() const { return
_condition_cache_hit_count; }
+ virtual int64_t condition_cache_hit_count() const { return
_condition_cache_hit_count; }
virtual std::string debug_string() const;
diff --git a/be/test/format_v2/column_mapper_test.cpp
b/be/test/format_v2/column_mapper_test.cpp
index 6bfc06139e4..d3428b6134a 100644
--- a/be/test/format_v2/column_mapper_test.cpp
+++ b/be/test/format_v2/column_mapper_test.cpp
@@ -2071,7 +2071,8 @@ TEST(ColumnMapperConstantTest,
PartitionDefaultAndVirtualColumnsUseDedicatedBran
{"dt", Field::create_field<TYPE_STRING>("2026-06-11")},
};
- TableColumnMapper mapper({.mode = TableColumnMappingMode::BY_NAME});
+ TableColumnMapper mapper(
+ {.mode = TableColumnMappingMode::BY_NAME,
.enable_row_lineage_virtual_columns = true});
ASSERT_TRUE(mapper.create_mapping(table_schema, partition_values,
{}).ok());
ASSERT_EQ(mapper.mappings().size(), 5);
@@ -2092,7 +2093,8 @@ TEST(ColumnMapperConstantTest,
PhysicalRowLineageFiltersStayFinalizeOnly) {
name_col("_last_updated_sequence_number", make_nullable(i64()),
2147483539),
};
- TableColumnMapper mapper({.mode = TableColumnMappingMode::BY_NAME});
+ TableColumnMapper mapper(
+ {.mode = TableColumnMappingMode::BY_NAME,
.enable_row_lineage_virtual_columns = true});
ASSERT_TRUE(mapper.create_mapping(table_schema, {}, file_schema).ok());
ASSERT_EQ(mapper.mappings().size(), 2);
@@ -2125,6 +2127,25 @@ TEST(ColumnMapperConstantTest,
PhysicalRowLineageFiltersStayFinalizeOnly) {
std::vector<int32_t>({2147483540, 2147483539}));
}
+TEST(ColumnMapperConstantTest, GenericByNameKeepsRowLineageNamesPhysical) {
+ const std::vector<ColumnDefinition> table_schema = {
+ name_col("_row_id", make_nullable(i64())),
+ name_col("_last_updated_sequence_number", make_nullable(i64())),
+ };
+ const std::vector<ColumnDefinition> file_schema = {
+ name_col("_row_id", make_nullable(i64()), 0),
+ name_col("_last_updated_sequence_number", make_nullable(i64()), 1),
+ };
+
+ TableColumnMapper mapper({.mode = TableColumnMappingMode::BY_NAME});
+ ASSERT_TRUE(mapper.create_mapping(table_schema, {}, file_schema).ok());
+ ASSERT_EQ(mapper.mappings().size(), 2);
+ EXPECT_EQ(mapper.mappings()[0].virtual_column_type,
TableVirtualColumnType::INVALID);
+ EXPECT_EQ(mapper.mappings()[0].filter_conversion,
FilterConversionType::COPY_DIRECTLY);
+ EXPECT_EQ(mapper.mappings()[1].virtual_column_type,
TableVirtualColumnType::INVALID);
+ EXPECT_EQ(mapper.mappings()[1].filter_conversion,
FilterConversionType::COPY_DIRECTLY);
+}
+
TEST(ColumnMapperConstantTest,
MissingRowLineageDefaultExprStillUsesVirtualMapping) {
auto id_column = field_id_col("id", 1, make_nullable(i32()));
auto row_id_column = field_id_col("renamed_row_id", 2147483540,
make_nullable(i64()));
@@ -2141,7 +2162,8 @@ TEST(ColumnMapperConstantTest,
MissingRowLineageDefaultExprStillUsesVirtualMappi
field_id_col("name", 2, make_nullable(str()), 1),
};
- TableColumnMapper mapper({.mode = TableColumnMappingMode::BY_FIELD_ID});
+ TableColumnMapper mapper({.mode = TableColumnMappingMode::BY_FIELD_ID,
+ .enable_row_lineage_virtual_columns = true});
ASSERT_TRUE(mapper.create_mapping(table_schema, {}, file_schema).ok());
ASSERT_EQ(mapper.mappings().size(), 3);
diff --git a/be/test/format_v2/jni/jdbc_reader_test.cpp
b/be/test/format_v2/jni/jdbc_reader_test.cpp
new file mode 100644
index 00000000000..5763a6deb00
--- /dev/null
+++ b/be/test/format_v2/jni/jdbc_reader_test.cpp
@@ -0,0 +1,49 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include "format_v2/jni/jdbc_reader.h"
+
+#include <gtest/gtest.h>
+
+#include "core/column/column_nullable.h"
+#include "core/column/column_string.h"
+#include "core/column/column_vector.h"
+
+namespace doris::format::jdbc {
+
+TEST(JdbcJniReaderTest, NonNullableSpecialTypeRejectsCastNull) {
+ auto data = ColumnString::create();
+ data->insert_default();
+ auto null_map = ColumnUInt8::create();
+ null_map->get_data().push_back(1);
+ auto result = ColumnNullable::create(std::move(data), std::move(null_map));
+
+ EXPECT_TRUE(validate_non_nullable_special_type_result(*result, 1)
+ .is<ErrorCode::DATA_QUALITY_ERROR>());
+}
+
+TEST(JdbcJniReaderTest, NonNullableSpecialTypeAcceptsSuccessfulCast) {
+ auto data = ColumnString::create();
+ data->insert_data("ok", 2);
+ auto null_map = ColumnUInt8::create();
+ null_map->get_data().push_back(0);
+ auto result = ColumnNullable::create(std::move(data), std::move(null_map));
+
+ EXPECT_TRUE(validate_non_nullable_special_type_result(*result, 1).ok());
+}
+
+} // namespace doris::format::jdbc
diff --git a/be/test/format_v2/json/json_reader_test.cpp
b/be/test/format_v2/json/json_reader_test.cpp
index 683a47dee29..e8d9aec40af 100644
--- a/be/test/format_v2/json/json_reader_test.cpp
+++ b/be/test/format_v2/json/json_reader_test.cpp
@@ -173,14 +173,21 @@ struct ReadResult {
bool eof = false;
size_t second_rows = 0;
bool second_eof = false;
+ size_t document_buffer_size = 0;
std::vector<ColumnDefinition> schema;
};
ReadResult read_once(const std::string& file_name, const std::string& content,
TFileScanRangeParams params, const
std::vector<SlotDescriptor*>& slots,
- const std::vector<int32_t>& requested_local_ids, bool
read_twice = false) {
+ const std::vector<int32_t>& requested_local_ids, bool
read_twice = false,
+ bool is_hive_table = false) {
const auto file_path = write_json_file(file_name, content);
auto range = file_range(file_path);
+ if (is_hive_table) {
+ TTableFormatFileDesc table_format;
+ table_format.__set_table_format_type("hive");
+ range.__set_table_format_params(std::move(table_format));
+ }
auto system_properties = std::make_shared<io::FileSystemProperties>();
system_properties->system_type = TFileType::FILE_LOCAL;
@@ -210,6 +217,7 @@ ReadResult read_once(const std::string& file_name, const
std::string& content,
result.block = make_block(result.schema, requested_local_ids);
result.status = reader.get_block(&result.block, &result.rows, &result.eof);
+ result.document_buffer_size = reader.TEST_document_buffer_size();
if (result.status.ok() && read_twice) {
auto eof_block = make_block(result.schema, requested_local_ids);
result.second_status =
@@ -318,6 +326,20 @@ TEST(JsonReaderTest,
ReadsRequestedColumnsInFileScanRequestOrder) {
ASSERT_TRUE(result.second_status.ok()) << result.second_status.to_string();
EXPECT_EQ(result.second_rows, 0);
EXPECT_TRUE(result.second_eof);
+ EXPECT_EQ(result.document_buffer_size, 0);
+}
+
+TEST(JsonReaderTest, HiveColumnLookupIsCaseInsensitiveWithoutNormalizedKeys) {
+ ObjectPool pool;
+ auto slots = build_slots(&pool);
+ auto result = read_once("hive_case.jsonl",
+ R"({"ID":7,"NaMe":"alice"})"
+ "\n",
+ json_scan_params(), slots, {0, 1}, false, true);
+ ASSERT_TRUE(result.status.ok()) << result.status.to_string();
+ ASSERT_EQ(result.rows, 1);
+ EXPECT_EQ(nullable_int_at(*result.block.get_by_position(0).column, 0), 7);
+ EXPECT_EQ(nullable_string_at(*result.block.get_by_position(1).column, 0),
"alice");
}
TEST(JsonReaderTest, ReadsSingleDocumentOuterArray) {
diff --git a/be/test/format_v2/parquet/native_decoder_test.cpp
b/be/test/format_v2/parquet/native_decoder_test.cpp
index 1cf6fdaade7..599c930f312 100644
--- a/be/test/format_v2/parquet/native_decoder_test.cpp
+++ b/be/test/format_v2/parquet/native_decoder_test.cpp
@@ -570,7 +570,8 @@ Status materialize_level_only_page(bool data_page_v2,
tparquet::Type::type physi
}
Status load_scripted_page(tparquet::PageHeader header, const
std::vector<uint8_t>& payload,
- tparquet::CompressionCodec::type codec, bool
preload_page_cache = false) {
+ tparquet::CompressionCodec::type codec, bool
preload_page_cache = false,
+ tparquet::Type::type physical_type =
tparquet::Type::INT32) {
std::vector<uint8_t> bytes;
ThriftSerializer serializer(/*compact=*/true, 128);
RETURN_IF_ERROR(serializer.serialize(&header, &bytes));
@@ -592,7 +593,7 @@ Status load_scripted_page(tparquet::PageHeader header,
const std::vector<uint8_t
MemoryBufferedReader reader(std::move(bytes));
tparquet::ColumnChunk chunk;
- chunk.meta_data.__set_type(tparquet::Type::INT32);
+ chunk.meta_data.__set_type(physical_type);
chunk.meta_data.__set_codec(codec);
chunk.meta_data.__set_num_values(1);
chunk.meta_data.__set_total_compressed_size(chunk_size);
@@ -603,8 +604,12 @@ Status load_scripted_page(tparquet::PageHeader header,
const std::vector<uint8_t
chunk.meta_data.__set_data_page_offset(0);
}
NativeFieldSchema field;
- field.physical_type = tparquet::Type::INT32;
- field.data_type = std::make_shared<DataTypeInt32>();
+ field.physical_type = physical_type;
+ if (physical_type == tparquet::Type::BYTE_ARRAY) {
+ field.data_type = std::make_shared<DataTypeString>();
+ } else {
+ field.data_type = std::make_shared<DataTypeInt32>();
+ }
field.repetition_level = 0;
field.definition_level = 0;
ParquetPageReadContext context(preload_page_cache, page_cache_file_key);
@@ -1279,6 +1284,80 @@ TEST(ParquetV2NativeDecoderTest,
DictionaryMaterializationUsesCacheAwareExecutio
verify_strategy(true, ParquetDictionaryMaterializationStrategy::INDICES,
0, 1);
}
+TEST(ParquetV2NativeDecoderTest,
DictionaryProbeMaterializesTypedValuesOnlyOnce) {
+ const std::array<int32_t, 2> dictionary {10, 20};
+ std::vector<uint8_t> dictionary_payload(sizeof(dictionary));
+ memcpy(dictionary_payload.data(), dictionary.data(),
dictionary_payload.size());
+ tparquet::PageHeader dictionary_header;
+ dictionary_header.type = tparquet::PageType::DICTIONARY_PAGE;
+ dictionary_header.__set_compressed_page_size(dictionary_payload.size());
+ dictionary_header.__set_uncompressed_page_size(dictionary_payload.size());
+ dictionary_header.__isset.dictionary_page_header = true;
+
dictionary_header.dictionary_page_header.__set_num_values(dictionary.size());
+
dictionary_header.dictionary_page_header.__set_encoding(tparquet::Encoding::PLAIN);
+ std::vector<uint8_t> bytes(1, 0);
+ const auto dictionary_page = serialize_page(dictionary_header,
dictionary_payload);
+ bytes.insert(bytes.end(), dictionary_page.begin(), dictionary_page.end());
+ const size_t data_page_offset = bytes.size();
+
+ faststring encoded_ids;
+ RleEncoder<uint32_t> encoder(&encoded_ids, 1);
+ for (const uint32_t id : {1U, 0U, 0U, 0U, 0U, 0U, 0U, 0U}) {
+ encoder.Put(id);
+ }
+ encoder.Flush();
+ std::vector<uint8_t> data_payload(encoded_ids.size() + 1);
+ data_payload[0] = 1;
+ memcpy(data_payload.data() + 1, encoded_ids.data(), encoded_ids.size());
+ tparquet::PageHeader data_header;
+ data_header.type = tparquet::PageType::DATA_PAGE;
+ data_header.__set_compressed_page_size(data_payload.size());
+ data_header.__set_uncompressed_page_size(data_payload.size());
+ data_header.__isset.data_page_header = true;
+ data_header.data_page_header.__set_num_values(2);
+
data_header.data_page_header.__set_encoding(tparquet::Encoding::RLE_DICTIONARY);
+
data_header.data_page_header.__set_definition_level_encoding(tparquet::Encoding::RLE);
+
data_header.data_page_header.__set_repetition_level_encoding(tparquet::Encoding::RLE);
+ const auto data_page = serialize_page(data_header, data_payload);
+ bytes.insert(bytes.end(), data_page.begin(), data_page.end());
+
+ tparquet::ColumnChunk chunk;
+ chunk.meta_data.__set_type(tparquet::Type::INT32);
+ chunk.meta_data.__set_codec(tparquet::CompressionCodec::UNCOMPRESSED);
+ chunk.meta_data.__set_num_values(2);
+ chunk.meta_data.__set_total_compressed_size(bytes.size() - 1);
+ chunk.meta_data.__set_dictionary_page_offset(1);
+ chunk.meta_data.__set_data_page_offset(data_page_offset);
+ NativeFieldSchema field;
+ field.physical_type = tparquet::Type::INT32;
+ field.data_type = std::make_shared<DataTypeInt32>();
+ field.parquet_schema.__set_type(tparquet::Type::INT32);
+
field.parquet_schema.__set_repetition_type(tparquet::FieldRepetitionType::REQUIRED);
+ auto file = std::make_shared<NativeDecoderMemoryFileReader>(bytes);
+ const auto row_ranges = ::doris::RowRanges::create_single(2);
+ ScalarColumnReader<false, false> reader(row_ranges, 2, chunk, nullptr,
nullptr, nullptr);
+ ASSERT_TRUE(reader.init(file, &field, bytes.size(), nullptr, "",
ParquetReaderCompat {}, true)
+ .ok());
+ auto dictionary_result = reader.dictionary_values(field.data_type);
+ ASSERT_TRUE(dictionary_result.has_value()) << dictionary_result.error();
+ EXPECT_EQ(reader.dictionary_materialization_count_for_test(), 1);
+
+ FilterMap filter;
+ ASSERT_TRUE(filter.init(nullptr, 2, false).ok());
+ ColumnPtr ids = ColumnInt32::create();
+ size_t rows = 0;
+ bool eof = false;
+ ASSERT_TRUE(reader.read_column_data(ids, field.data_type, nullptr, filter,
2, &rows, &eof, true)
+ .ok());
+ ASSERT_EQ(rows, 2);
+ auto matched_values = reader.materialize_dictionary_values(
+ &assert_cast<const ColumnInt32&>(*ids), field.data_type);
+ ASSERT_TRUE(matched_values.has_value()) << matched_values.error();
+ EXPECT_EQ(reader.dictionary_materialization_count_for_test(), 1);
+ EXPECT_EQ(assert_cast<const ColumnInt32&>(**matched_values).get_data(),
+ (ColumnInt32::Container {20, 10}));
+}
+
TEST(ParquetV2NativeDecoderTest,
DictionaryRepeatedRunsGatherDirectlyIntoDestination) {
const std::array<int32_t, 2> dictionary_values {10, 20};
auto dictionary = make_unique_buffer<uint8_t>(sizeof(dictionary_values));
@@ -2009,6 +2088,24 @@ TEST(ParquetV2NativeDecoderTest,
DeltaEncodingsExposeValuesAfterSkip) {
}
}
+TEST(ParquetV2NativeDecoderTest,
DeltaBinaryPackedRejectsNonIntegralBlockGeometry) {
+ std::vector<uint8_t> encoded(32);
+ uint8_t* cursor = encoded.data();
+ cursor = encode_varint32(cursor, 3200);
+ cursor = encode_varint32(cursor, 33);
+ cursor = encode_varint32(cursor, 1);
+ cursor = encode_varint32(cursor, 0);
+ encoded.resize(cursor - encoded.data());
+
+ std::unique_ptr<Decoder> decoder;
+ ASSERT_TRUE(Decoder::get_decoder(tparquet::Type::INT32,
tparquet::Encoding::DELTA_BINARY_PACKED,
+ decoder)
+ .ok());
+ decoder->set_expected_values(1);
+ Slice slice(encoded.data(), encoded.size());
+ EXPECT_FALSE(decoder->set_data(&slice).ok());
+}
+
TEST(ParquetV2NativeDecoderTest, SparseStatefulEncodingsBatchDecodeAndCompact)
{
const ParquetSelection selection {
.total_values = 3,
@@ -2699,6 +2796,73 @@ TEST(ParquetV2NativeDecoderTest,
DecoderOwnedHighWaterScratchIsReleased) {
EXPECT_LE(decoder->retained_scratch_bytes(), 64UL << 10);
}
+TEST(ParquetV2NativeDecoderTest,
DecompressionScratchStaysActiveUntilPageExhaustion) {
+ BlockCompressionCodec* codec = nullptr;
+
ASSERT_TRUE(get_block_compression_codec(tparquet::CompressionCodec::SNAPPY,
&codec).ok());
+ auto compressed_page = [&](uint32_t value_count) {
+ std::vector<uint8_t> encoded_levels(8);
+ uint8_t* level_end = encode_varint32(encoded_levels.data(),
value_count << 1);
+ *level_end++ = 1;
+ encoded_levels.resize(level_end - encoded_levels.data());
+ std::vector<uint8_t> payload(sizeof(uint32_t));
+ encode_fixed32_le(payload.data(), encoded_levels.size());
+ payload.insert(payload.end(), encoded_levels.begin(),
encoded_levels.end());
+ payload.resize(payload.size() + static_cast<size_t>(value_count) *
sizeof(int32_t));
+
+ faststring compressed;
+ DORIS_CHECK(codec->compress(Slice(payload.data(), payload.size()),
&compressed).ok());
+ tparquet::PageHeader header;
+ header.type = tparquet::PageType::DATA_PAGE;
+ header.__set_compressed_page_size(compressed.size());
+ header.__set_uncompressed_page_size(payload.size());
+ header.__isset.data_page_header = true;
+ header.data_page_header.__set_num_values(value_count);
+ header.data_page_header.__set_encoding(tparquet::Encoding::PLAIN);
+
header.data_page_header.__set_definition_level_encoding(tparquet::Encoding::RLE);
+
header.data_page_header.__set_repetition_level_encoding(tparquet::Encoding::RLE);
+ return serialize_page(header, std::vector<uint8_t>(compressed.data(),
+ compressed.data() +
compressed.size()));
+ };
+
+ constexpr uint32_t LARGE_VALUE_COUNT = 1U << 20;
+ auto bytes = compressed_page(LARGE_VALUE_COUNT);
+ const auto ordinary_page = compressed_page(1);
+ bytes.insert(bytes.end(), ordinary_page.begin(), ordinary_page.end());
+ MemoryBufferedReader stream(bytes);
+ tparquet::ColumnChunk chunk;
+ chunk.meta_data.__set_type(tparquet::Type::INT32);
+ chunk.meta_data.__set_codec(tparquet::CompressionCodec::SNAPPY);
+ chunk.meta_data.__set_num_values(static_cast<int64_t>(LARGE_VALUE_COUNT) +
1);
+ chunk.meta_data.__set_total_compressed_size(bytes.size());
+ chunk.meta_data.__set_data_page_offset(0);
+ NativeFieldSchema field;
+ field.physical_type = tparquet::Type::INT32;
+ field.definition_level = 1;
+
field.parquet_schema.__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL);
+ ParquetPageReadContext context(false, "");
+ ColumnChunkReader<false, false> reader(&stream, &chunk, &field, nullptr,
+
static_cast<size_t>(LARGE_VALUE_COUNT) + 1, nullptr,
+ context);
+ ASSERT_TRUE(reader.init().ok());
+ ASSERT_TRUE(reader.load_page_data().ok());
+ ASSERT_GT(reader.active_decoder_scratch_bytes(), 1UL << 20);
+ ASSERT_TRUE(reader.skip_values(LARGE_VALUE_COUNT).ok());
+ ASSERT_TRUE(reader.next_page().ok());
+ ASSERT_TRUE(reader.parse_page_header().ok());
+ ASSERT_TRUE(reader.load_page_data().ok());
+ EXPECT_LT(reader.active_decoder_scratch_bytes(), 64UL << 10);
+ const size_t retained_while_active =
reader.retained_decoder_scratch_bytes();
+ ASSERT_GT(retained_while_active, 1UL << 20);
+ reader.release_decoder_scratch(64UL << 10);
+ EXPECT_EQ(reader.retained_decoder_scratch_bytes(), retained_while_active);
+
+ ASSERT_TRUE(reader.skip_values(1).ok());
+ ASSERT_TRUE(reader.next_page().ok());
+ // A release requested while a decoder points into the buffer must execute
at the next safe
+ // page boundary; otherwise periodic probes can never age out a previous
large-page capacity.
+ EXPECT_LE(reader.retained_decoder_scratch_bytes(), 64UL << 10);
+}
+
TEST(ParquetV2NativeDecoderTest,
DeltaByteArrayScratchReleasePreservesPrefixState) {
const std::vector<std::string> values {std::string(4096, 'x'),
"shared-prefix-a",
"shared-prefix-b",
"shared-prefix-c"};
@@ -3386,6 +3550,215 @@ TEST(ParquetV2NativeDecoderTest,
UncompressedDictionaryRequiresEqualPhysicalAndL
.is<ErrorCode::CORRUPTION>());
}
+TEST(ParquetV2NativeDecoderTest,
EmptyDictionaryRejectsDeclaredPayloadBeforeAllocation) {
+ tparquet::PageHeader header;
+ header.__set_uncompressed_page_size(std::numeric_limits<int32_t>::max());
+ header.__isset.dictionary_page_header = true;
+ header.dictionary_page_header.__set_num_values(0);
+
EXPECT_TRUE(validate_dictionary_page_size(header).is<ErrorCode::CORRUPTION>());
+}
+
+TEST(ParquetV2NativeDecoderTest,
NonemptyFixedWidthDictionaryRejectsExtentBeforeAllocation) {
+ BlockCompressionCodec* codec = nullptr;
+
ASSERT_TRUE(get_block_compression_codec(tparquet::CompressionCodec::SNAPPY,
&codec).ok());
+ const std::array<uint8_t, sizeof(int32_t)> value {};
+ faststring compressed;
+ ASSERT_TRUE(codec->compress(Slice(value.data(), value.size()),
&compressed).ok());
+
+ tparquet::PageHeader header;
+ header.type = tparquet::PageType::DICTIONARY_PAGE;
+ header.__set_compressed_page_size(compressed.size());
+ header.__set_uncompressed_page_size(8 << 20);
+ header.__isset.dictionary_page_header = true;
+ header.dictionary_page_header.__set_num_values(1);
+ header.dictionary_page_header.__set_encoding(tparquet::Encoding::PLAIN);
+ EXPECT_TRUE(validate_dictionary_page_size(header,
sizeof(int32_t)).is<ErrorCode::CORRUPTION>());
+
+ const std::vector<uint8_t> payload(compressed.data(), compressed.data() +
compressed.size());
+ EXPECT_TRUE(load_scripted_page(header, payload,
tparquet::CompressionCodec::SNAPPY)
+ .is<ErrorCode::CORRUPTION>());
+}
+
+TEST(ParquetV2NativeDecoderTest,
RequiredFixedWidthPageRejectsImpossibleExtentBeforeAllocation) {
+ for (const auto encoding : {tparquet::Encoding::PLAIN,
tparquet::Encoding::BYTE_STREAM_SPLIT}) {
+ tparquet::PageHeader header;
+
header.__set_uncompressed_page_size(std::numeric_limits<int32_t>::max());
+ header.__isset.data_page_header = true;
+ header.data_page_header.__set_num_values(1);
+ header.data_page_header.__set_encoding(encoding);
+ EXPECT_TRUE(validate_fixed_width_page_size(header, sizeof(int32_t), 0,
0)
+ .is<ErrorCode::CORRUPTION>());
+ header.__set_uncompressed_page_size(sizeof(int32_t));
+ EXPECT_TRUE(validate_fixed_width_page_size(header, sizeof(int32_t), 0,
0).ok());
+ }
+}
+
+TEST(ParquetV2NativeDecoderTest,
OptionalV2FixedWidthPageRejectsExtentBeforeAllocation) {
+ BlockCompressionCodec* codec = nullptr;
+
ASSERT_TRUE(get_block_compression_codec(tparquet::CompressionCodec::SNAPPY,
&codec).ok());
+ const std::array<uint8_t, sizeof(int32_t)> value {};
+ faststring compressed;
+ ASSERT_TRUE(codec->compress(Slice(value.data(), value.size()),
&compressed).ok());
+ const std::vector<uint8_t> levels {2, 1};
+
+ for (const auto encoding : {tparquet::Encoding::PLAIN,
tparquet::Encoding::BYTE_STREAM_SPLIT}) {
+ tparquet::PageHeader header;
+ header.type = tparquet::PageType::DATA_PAGE_V2;
+ header.__set_compressed_page_size(levels.size() + compressed.size());
+ header.__set_uncompressed_page_size((8 << 20) + levels.size());
+ header.__isset.data_page_header_v2 = true;
+ header.data_page_header_v2.__set_num_values(1);
+ header.data_page_header_v2.__set_num_rows(1);
+ header.data_page_header_v2.__set_num_nulls(0);
+ header.data_page_header_v2.__set_encoding(encoding);
+ header.data_page_header_v2.__set_repetition_levels_byte_length(0);
+
header.data_page_header_v2.__set_definition_levels_byte_length(levels.size());
+ header.data_page_header_v2.__set_is_compressed(true);
+ std::vector<uint8_t> payload = levels;
+ payload.insert(payload.end(), compressed.data(), compressed.data() +
compressed.size());
+
+ auto bytes = serialize_page(header, payload);
+ MemoryBufferedReader stream(bytes);
+ tparquet::ColumnChunk chunk;
+ chunk.meta_data.__set_type(tparquet::Type::INT32);
+ chunk.meta_data.__set_codec(tparquet::CompressionCodec::SNAPPY);
+ chunk.meta_data.__set_num_values(1);
+ chunk.meta_data.__set_total_compressed_size(bytes.size());
+ chunk.meta_data.__set_data_page_offset(0);
+ NativeFieldSchema field;
+ field.physical_type = tparquet::Type::INT32;
+ field.definition_level = 1;
+ field.parquet_schema.__set_type(tparquet::Type::INT32);
+
field.parquet_schema.__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL);
+ ParquetPageReadContext context(false, "");
+ ColumnChunkReader<false, false> reader(&stream, &chunk, &field,
nullptr, 1, nullptr,
+ context);
+ ASSERT_TRUE(reader.init().ok());
+ EXPECT_TRUE(reader.load_page_data().is<ErrorCode::CORRUPTION>());
+ EXPECT_LT(reader.retained_decoder_scratch_bytes(), 64UL << 10);
+ }
+}
+
+TEST(ParquetV2NativeDecoderTest,
RepeatedV2FixedWidthPageRejectsExtentBeforeAllocation) {
+ BlockCompressionCodec* codec = nullptr;
+
ASSERT_TRUE(get_block_compression_codec(tparquet::CompressionCodec::SNAPPY,
&codec).ok());
+ const std::array<uint8_t, sizeof(int32_t)> value {};
+ faststring compressed;
+ ASSERT_TRUE(codec->compress(Slice(value.data(), value.size()),
&compressed).ok());
+ const std::vector<uint8_t> levels {2, 0, 2, 1};
+
+ tparquet::PageHeader header;
+ header.type = tparquet::PageType::DATA_PAGE_V2;
+ header.__set_compressed_page_size(levels.size() + compressed.size());
+ header.__set_uncompressed_page_size((8 << 20) + levels.size());
+ header.__isset.data_page_header_v2 = true;
+ header.data_page_header_v2.__set_num_values(1);
+ header.data_page_header_v2.__set_num_rows(1);
+ header.data_page_header_v2.__set_num_nulls(0);
+ header.data_page_header_v2.__set_encoding(tparquet::Encoding::PLAIN);
+ header.data_page_header_v2.__set_repetition_levels_byte_length(2);
+ header.data_page_header_v2.__set_definition_levels_byte_length(2);
+ header.data_page_header_v2.__set_is_compressed(true);
+ std::vector<uint8_t> payload = levels;
+ payload.insert(payload.end(), compressed.data(), compressed.data() +
compressed.size());
+
+ auto bytes = serialize_page(header, payload);
+ MemoryBufferedReader stream(bytes);
+ tparquet::ColumnChunk chunk;
+ chunk.meta_data.__set_type(tparquet::Type::INT32);
+ chunk.meta_data.__set_codec(tparquet::CompressionCodec::SNAPPY);
+ chunk.meta_data.__set_num_values(1);
+ chunk.meta_data.__set_total_compressed_size(bytes.size());
+ chunk.meta_data.__set_data_page_offset(0);
+ NativeFieldSchema field;
+ field.physical_type = tparquet::Type::INT32;
+ field.repetition_level = 1;
+ field.definition_level = 1;
+ field.parquet_schema.__set_type(tparquet::Type::INT32);
+
field.parquet_schema.__set_repetition_type(tparquet::FieldRepetitionType::REPEATED);
+ ParquetPageReadContext context(false, "");
+ ColumnChunkReader<true, false> reader(&stream, &chunk, &field, nullptr, 1,
nullptr, context);
+ ASSERT_TRUE(reader.init().ok());
+ EXPECT_TRUE(reader.load_page_data().is<ErrorCode::CORRUPTION>());
+ EXPECT_LT(reader.retained_decoder_scratch_bytes(), 64UL << 10);
+}
+
+TEST(ParquetV2NativeDecoderTest,
VariableWidthDataPagePreflightsCompressedExtent) {
+ BlockCompressionCodec* codec = nullptr;
+
ASSERT_TRUE(get_block_compression_codec(tparquet::CompressionCodec::SNAPPY,
&codec).ok());
+ const std::vector<uint8_t> value {1, 0, 0, 0, 'x'};
+ faststring compressed;
+ ASSERT_TRUE(codec->compress(Slice(value.data(), value.size()),
&compressed).ok());
+
+ tparquet::PageHeader header;
+ header.type = tparquet::PageType::DATA_PAGE;
+ header.__set_compressed_page_size(compressed.size());
+ header.__set_uncompressed_page_size(8 << 20);
+ header.__isset.data_page_header = true;
+ header.data_page_header.__set_num_values(1);
+ header.data_page_header.__set_encoding(tparquet::Encoding::PLAIN);
+
header.data_page_header.__set_repetition_level_encoding(tparquet::Encoding::RLE);
+
header.data_page_header.__set_definition_level_encoding(tparquet::Encoding::RLE);
+ const std::vector<uint8_t> payload(compressed.data(), compressed.data() +
compressed.size());
+
+ auto bytes = serialize_page(header, payload);
+ MemoryBufferedReader stream(bytes);
+ tparquet::ColumnChunk chunk;
+ chunk.meta_data.__set_type(tparquet::Type::BYTE_ARRAY);
+ chunk.meta_data.__set_codec(tparquet::CompressionCodec::SNAPPY);
+ chunk.meta_data.__set_num_values(1);
+ chunk.meta_data.__set_total_compressed_size(bytes.size());
+ chunk.meta_data.__set_data_page_offset(0);
+ NativeFieldSchema field;
+ field.physical_type = tparquet::Type::BYTE_ARRAY;
+ field.definition_level = 1;
+ field.parquet_schema.__set_type(tparquet::Type::BYTE_ARRAY);
+
field.parquet_schema.__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL);
+ ParquetPageReadContext context(false, "");
+ ColumnChunkReader<false, false> reader(&stream, &chunk, &field, nullptr,
1, nullptr, context);
+ ASSERT_TRUE(reader.init().ok());
+ EXPECT_TRUE(reader.load_page_data().is<ErrorCode::CORRUPTION>());
+ EXPECT_LT(reader.retained_decoder_scratch_bytes(), 64UL << 10);
+ EXPECT_TRUE(load_scripted_page(header, payload,
tparquet::CompressionCodec::SNAPPY, true,
+ tparquet::Type::BYTE_ARRAY)
+ .is<ErrorCode::CORRUPTION>());
+}
+
+TEST(ParquetV2NativeDecoderTest,
VariableWidthDictionaryPreflightsCompressedExtent) {
+ BlockCompressionCodec* codec = nullptr;
+
ASSERT_TRUE(get_block_compression_codec(tparquet::CompressionCodec::SNAPPY,
&codec).ok());
+ const std::vector<uint8_t> value {1, 0, 0, 0, 'x'};
+ faststring compressed;
+ ASSERT_TRUE(codec->compress(Slice(value.data(), value.size()),
&compressed).ok());
+ const Slice payload(compressed.data(), compressed.size());
+
+ EXPECT_TRUE(
+ validate_compressed_page_size(tparquet::CompressionCodec::SNAPPY,
payload, value.size())
+ .ok());
+
EXPECT_TRUE(validate_compressed_page_size(tparquet::CompressionCodec::SNAPPY,
payload, 8 << 20)
+ .is<ErrorCode::CORRUPTION>());
+
+ tparquet::PageHeader header;
+ header.type = tparquet::PageType::DICTIONARY_PAGE;
+ header.__set_compressed_page_size(compressed.size());
+ header.__set_uncompressed_page_size(8 << 20);
+ header.__isset.dictionary_page_header = true;
+ header.dictionary_page_header.__set_num_values(1);
+ header.dictionary_page_header.__set_encoding(tparquet::Encoding::PLAIN);
+ EXPECT_TRUE(
+ load_scripted_page(
+ header,
+ std::vector<uint8_t>(compressed.data(), compressed.data()
+ compressed.size()),
+ tparquet::CompressionCodec::SNAPPY, false,
tparquet::Type::BYTE_ARRAY)
+ .is<ErrorCode::CORRUPTION>());
+ EXPECT_TRUE(
+ load_scripted_page(
+ header,
+ std::vector<uint8_t>(compressed.data(), compressed.data()
+ compressed.size()),
+ tparquet::CompressionCodec::SNAPPY, true,
tparquet::Type::BYTE_ARRAY)
+ .is<ErrorCode::CORRUPTION>());
+}
+
TEST(ParquetV2NativeDecoderTest,
UncompressedDataPagesRequireEqualPhysicalAndLogicalSizes) {
for (const auto page_type : {tparquet::PageType::DATA_PAGE,
tparquet::PageType::DATA_PAGE_V2}) {
tparquet::PageHeader header;
diff --git a/be/test/format_v2/parquet/parquet_scan_test.cpp
b/be/test/format_v2/parquet/parquet_scan_test.cpp
index cc1d28238cb..9bdb0c7562f 100644
--- a/be/test/format_v2/parquet/parquet_scan_test.cpp
+++ b/be/test/format_v2/parquet/parquet_scan_test.cpp
@@ -928,6 +928,14 @@ TEST(ParquetScanAdaptivePredicateTest,
SamplesWarmupThenAtLowFrequency) {
EXPECT_TRUE(should_sample_adaptive_predicate(9, 32));
}
+TEST(ParquetScanDeleteConjunctTest, RejectsInputColumnAsEphemeralResult) {
+
EXPECT_TRUE(format::parquet::detail::validate_ephemeral_expr_result_column(2,
0, 2)
+ .is<ErrorCode::INTERNAL_ERROR>());
+
EXPECT_TRUE(format::parquet::detail::validate_ephemeral_expr_result_column(2,
2, 3).ok());
+
EXPECT_TRUE(format::parquet::detail::validate_ephemeral_expr_result_column(2,
3, 3)
+ .is<ErrorCode::INTERNAL_ERROR>());
+}
+
TEST(ParquetScanAdaptivePredicateTest,
ThrowingNestedFunctionDisablesSelectedRowReordering) {
using format::parquet::detail::AdaptivePredicateStats;
std::unordered_map<size_t, AdaptivePredicateStats> first_batch_stats;
@@ -1630,6 +1638,48 @@ TEST_F(ParquetScanTest,
ProjectedDeltaBinaryPackedUsesFixedWidthFilterAndProject
EXPECT_EQ(counter_value(profile, "PredicateCompactionBytes"), 0);
}
+TEST_F(ParquetScanTest, PlainPredicateDirectPathCrossesScratchProbeCadence) {
+ constexpr int64_t ROWS = 32;
+ std::vector<int32_t> ids(ROWS);
+ std::vector<int32_t> scores(ROWS);
+ std::iota(ids.begin(), ids.end(), 1);
+ std::iota(scores.begin(), scores.end(), 10);
+ auto schema = arrow::schema({
+ arrow::field("id", arrow::int32(), false),
+ arrow::field("score", arrow::int32(), false),
+ });
+ auto table = arrow::Table::Make(schema, {build_int32_array(ids),
build_int32_array(scores)});
+ write_table(_file_path, table, ROWS, false, false, false);
+ RuntimeProfile profile("profile");
+ auto reader = create_reader(0, -1, &profile);
+ reader->set_batch_size(1);
+ RuntimeState state {TQueryOptions(), TQueryGlobals()};
+ ASSERT_TRUE(reader->init(&state).ok());
+
+ std::vector<format::ColumnDefinition> file_schema;
+ ASSERT_TRUE(reader->get_schema(&file_schema).ok());
+ auto request = std::make_shared<format::FileScanRequest>();
+ format::FileScanRequestBuilder request_builder(request.get());
+
ASSERT_TRUE(request_builder.add_predicate_column(format::LocalColumnId(0)).ok());
+
ASSERT_TRUE(request_builder.add_non_predicate_column(format::LocalColumnId(1)).ok());
+ request->predicate_only_columns.push_back(format::LocalColumnId(0));
+ request->conjuncts.push_back(create_int32_function_conjunct(0, "gt",
TExprOpcode::GT, 0));
+ ASSERT_TRUE(reader->open(request).ok());
+
+ size_t total_rows = 0;
+ bool eof = false;
+ while (!eof) {
+ Block block = build_file_block(file_schema);
+ size_t rows = 0;
+ ASSERT_TRUE(reader->get_block(&block, &rows, &eof).ok());
+ total_rows += rows;
+ }
+ EXPECT_EQ(total_rows, ROWS);
+ // Direct predicate evaluation must survive multiple 16-batch scratch
probes; otherwise this
+ // path can retain a previous outlier for the whole row group without ever
aging its capacity.
+ EXPECT_EQ(counter_value(profile, "FixedWidthPredicateDirectBatches"),
ROWS);
+}
+
TEST_F(ParquetScanTest, PredicateOnlyUint32FallsBackBeforeRawPlainDecode) {
write_uint32_pair_parquet_file(_file_path);
RuntimeProfile profile("profile");
diff --git a/be/test/format_v2/parquet/parquet_schema_test.cpp
b/be/test/format_v2/parquet/parquet_schema_test.cpp
index d116f84fbaf..74781c0e532 100644
--- a/be/test/format_v2/parquet/parquet_schema_test.cpp
+++ b/be/test/format_v2/parquet/parquet_schema_test.cpp
@@ -26,6 +26,7 @@
#include "core/data_type/data_type_array.h"
#include "core/data_type/data_type_map.h"
#include "core/data_type/data_type_nullable.h"
+#include "core/data_type/data_type_number.h"
#include "core/data_type/data_type_struct.h"
#include "core/data_type/primitive_type.h"
#include "format_v2/parquet/native_schema_desc.h"
@@ -606,4 +607,115 @@ TEST(ParquetSchemaTest,
NativeMetadataRejectsRowGroupChunkCardinalityAndMissingM
}
}
+TEST(ParquetSchemaTest,
NativeProjectionUsesResolvedIdentityBeforeCaseInsensitiveFallback) {
+ tparquet::SchemaElement root;
+ root.__set_name("schema");
+ root.__set_num_children(1);
+ tparquet::SchemaElement group;
+ group.__set_name("s");
+ group.__set_num_children(3);
+ group.__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL);
+ tparquet::SchemaElement upper;
+ upper.__set_name("Value");
+ upper.__set_type(tparquet::Type::INT32);
+ upper.__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL);
+ upper.__set_field_id(1);
+ tparquet::SchemaElement lower = upper;
+ lower.__set_name("value");
+ lower.__set_field_id(2);
+ tparquet::SchemaElement other = upper;
+ other.__set_name("other");
+ other.__set_field_id(3);
+
+ NativeFieldDescriptor descriptor;
+ ASSERT_TRUE(descriptor.parse_from_thrift({root, group, upper, lower,
other}).ok());
+ descriptor.assign_ids();
+ std::vector<std::unique_ptr<ParquetColumnSchema>> fields;
+ ASSERT_TRUE(build_parquet_column_schema(descriptor, &fields).ok());
+
+ const auto int_type = make_nullable(std::make_shared<DataTypeInt32>());
+ std::shared_ptr<NativeSchemaNode> mapping;
+ auto other_only = make_nullable(
+ std::make_shared<DataTypeStruct>(DataTypes {int_type}, Strings
{"other"}));
+ ASSERT_TRUE(build_native_schema_node(other_only, *fields[0],
&mapping).ok());
+ EXPECT_TRUE(mapping->has_child("other"));
+
+ auto exact_case = make_nullable(
+ std::make_shared<DataTypeStruct>(DataTypes {int_type}, Strings
{"Value"}));
+ ASSERT_TRUE(build_native_schema_node(exact_case, *fields[0],
&mapping).ok());
+ EXPECT_EQ(mapping->file_child_name("Value"), "Value");
+
+ auto ambiguous_fallback = make_nullable(
+ std::make_shared<DataTypeStruct>(DataTypes {int_type}, Strings
{"VALUE"}));
+ const auto status = build_native_schema_node(ambiguous_fallback,
*fields[0], &mapping);
+ EXPECT_TRUE(status.is<ErrorCode::CORRUPTION>()) << status;
+}
+
+TEST(ParquetSchemaTest, NativeSetMapKeyValueWrapperRemainsSingleListElement) {
+ tparquet::SchemaElement root;
+ root.__set_name("schema");
+ root.__set_num_children(1);
+ tparquet::SchemaElement set;
+ set.__set_name("tags");
+ set.__set_num_children(1);
+ set.__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL);
+ set.__set_converted_type(tparquet::ConvertedType::MAP);
+ tparquet::SchemaElement wrapper;
+ wrapper.__set_name("key_value");
+ wrapper.__set_num_children(1);
+ wrapper.__set_repetition_type(tparquet::FieldRepetitionType::REPEATED);
+ wrapper.__set_converted_type(tparquet::ConvertedType::MAP_KEY_VALUE);
+ tparquet::SchemaElement key;
+ key.__set_name("key");
+ key.__set_type(tparquet::Type::INT32);
+ key.__set_repetition_type(tparquet::FieldRepetitionType::REQUIRED);
+
+ NativeFieldDescriptor descriptor;
+ ASSERT_TRUE(descriptor.parse_from_thrift({root, set, wrapper, key}).ok());
+ const auto* column = descriptor.get_column(0);
+ EXPECT_EQ(remove_nullable(column->data_type)->get_primitive_type(),
TYPE_ARRAY);
+ ASSERT_EQ(column->children.size(), 1);
+
EXPECT_EQ(remove_nullable(column->children[0].data_type)->get_primitive_type(),
TYPE_INT);
+}
+
+TEST(ParquetSchemaTest, NativeListPreservesLegacyMapKeyValueElement) {
+ tparquet::SchemaElement root;
+ root.__set_name("schema");
+ root.__set_num_children(1);
+ tparquet::SchemaElement list;
+ list.__set_name("entries");
+ list.__set_num_children(1);
+ list.__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL);
+ list.__set_converted_type(tparquet::ConvertedType::LIST);
+ tparquet::SchemaElement element;
+ element.__set_name("element");
+ element.__set_num_children(1);
+ element.__set_repetition_type(tparquet::FieldRepetitionType::REPEATED);
+ element.__set_converted_type(tparquet::ConvertedType::MAP_KEY_VALUE);
+ tparquet::SchemaElement map;
+ map.__set_name("map");
+ map.__set_num_children(2);
+ map.__set_repetition_type(tparquet::FieldRepetitionType::REPEATED);
+ tparquet::SchemaElement key;
+ key.__set_name("key");
+ key.__set_type(tparquet::Type::INT32);
+ key.__set_repetition_type(tparquet::FieldRepetitionType::REQUIRED);
+ tparquet::SchemaElement value;
+ value.__set_name("value");
+ value.__set_type(tparquet::Type::BYTE_ARRAY);
+ value.__set_repetition_type(tparquet::FieldRepetitionType::OPTIONAL);
+
+ NativeFieldDescriptor descriptor;
+ ASSERT_TRUE(descriptor.parse_from_thrift({root, list, element, map, key,
value}).ok());
+ const auto* column = descriptor.get_column(0);
+ EXPECT_EQ(remove_nullable(column->data_type)->get_primitive_type(),
TYPE_ARRAY);
+ ASSERT_EQ(column->children.size(), 1);
+
EXPECT_EQ(remove_nullable(column->children[0].data_type)->get_primitive_type(),
TYPE_MAP);
+ ASSERT_EQ(column->children[0].children.size(), 2);
+
EXPECT_EQ(remove_nullable(column->children[0].children[0].data_type)->get_primitive_type(),
+ TYPE_INT);
+
EXPECT_EQ(remove_nullable(column->children[0].children[1].data_type)->get_primitive_type(),
+ TYPE_STRING);
+}
+
} // namespace doris::format::parquet
diff --git a/be/test/format_v2/table/hudi_reader_test.cpp
b/be/test/format_v2/table/hudi_reader_test.cpp
index b8dd97e57d5..f5bd70292e5 100644
--- a/be/test/format_v2/table/hudi_reader_test.cpp
+++ b/be/test/format_v2/table/hudi_reader_test.cpp
@@ -409,6 +409,13 @@ TEST(HudiHybridReaderTest,
ScannerStatefulResidualSurvivesNativeJniNativeSwitch)
EXPECT_EQ(observed_invocations, std::vector<int>({0, 1, 2}));
}
+TEST(HudiHybridReaderTest, AggregatesConditionCacheHitsFromBothChildren) {
+ hudi::HudiHybridReader reader;
+ reader.TEST_install_batch_size_children();
+ reader.TEST_set_child_condition_cache_hits(2, 7);
+ EXPECT_EQ(reader.condition_cache_hit_count(), 9);
+}
+
TEST(HudiHybridReaderTest,
NativeCountStarReportsMetadataRowsThroughHybridReader) {
const auto test_dir =
std::filesystem::temp_directory_path() /
"doris_hudi_hybrid_count_star_test";
diff --git a/be/test/format_v2/table/paimon_reader_test.cpp
b/be/test/format_v2/table/paimon_reader_test.cpp
index 2434d748dfa..77466439c6b 100644
--- a/be/test/format_v2/table/paimon_reader_test.cpp
+++ b/be/test/format_v2/table/paimon_reader_test.cpp
@@ -840,6 +840,13 @@ TEST(PaimonHybridReaderTest,
ScannerStatefulResidualSurvivesNativeJniNativeSwitc
EXPECT_EQ(observed_invocations, std::vector<int>({0, 1, 2}));
}
+TEST(PaimonHybridReaderTest, AggregatesConditionCacheHitsFromBothChildren) {
+ paimon::PaimonHybridReader reader;
+ reader.TEST_install_batch_size_children();
+ reader.TEST_set_child_condition_cache_hits(3, 5);
+ EXPECT_EQ(reader.condition_cache_hit_count(), 8);
+}
+
TEST(PaimonHybridReaderTest,
NativeCountColumnReportsMetadataRowsThroughHybridReader) {
const auto test_dir =
std::filesystem::temp_directory_path() /
"doris_paimon_hybrid_count_column_test";
diff --git a/be/test/format_v2/table/remote_doris_reader_test.cpp
b/be/test/format_v2/table/remote_doris_reader_test.cpp
index b17f82f505c..a8affbde311 100644
--- a/be/test/format_v2/table/remote_doris_reader_test.cpp
+++ b/be/test/format_v2/table/remote_doris_reader_test.cpp
@@ -18,11 +18,17 @@
#include "format_v2/table/remote_doris_reader.h"
#include <arrow/api.h>
+#include <arrow/flight/server.h>
#include <gtest/gtest.h>
#include <algorithm>
+#include <chrono>
+#include <condition_variable>
+#include <future>
#include <memory>
+#include <mutex>
#include <string>
+#include <thread>
#include <tuple>
#include <utility>
#include <vector>
@@ -48,6 +54,7 @@
#include "runtime/runtime_profile.h"
#include "runtime/runtime_state.h"
#include "testutil/desc_tbl_builder.h"
+#include "testutil/mock/mock_runtime_state.h"
namespace doris::format::remote_doris {
namespace {
@@ -95,6 +102,122 @@ TFileRangeDesc remote_doris_range() {
return range;
}
+class BlockingRecordBatchReader final : public arrow::RecordBatchReader {
+public:
+ std::shared_ptr<arrow::Schema> schema() const override {
+ return arrow::schema({arrow::field("id", arrow::int32())});
+ }
+
+ arrow::Status ReadNext(std::shared_ptr<arrow::RecordBatch>* batch)
override {
+ {
+ std::lock_guard lock(_mutex);
+ _entered = true;
+ }
+ _cv.notify_all();
+ std::unique_lock lock(_mutex);
+ _cv.wait(lock, [this] { return _released; });
+ *batch = nullptr;
+ return arrow::Status::OK();
+ }
+
+ bool wait_until_entered(std::chrono::milliseconds timeout) {
+ std::unique_lock lock(_mutex);
+ return _cv.wait_for(lock, timeout, [this] { return _entered; });
+ }
+
+ void release() {
+ {
+ std::lock_guard lock(_mutex);
+ _released = true;
+ }
+ _cv.notify_all();
+ }
+
+private:
+ std::mutex _mutex;
+ std::condition_variable _cv;
+ bool _entered = false;
+ bool _released = false;
+};
+
+class BlockingFlightServer final : public arrow::flight::FlightServerBase {
+public:
+ enum class Mode { DO_GET, NEXT };
+
+ explicit BlockingFlightServer(Mode mode)
+ : _mode(mode),
_batch_reader(std::make_shared<BlockingRecordBatchReader>()) {}
+
+ arrow::Status start() {
+ auto location = arrow::flight::Location::ForGrpcTcp("localhost", 0);
+ if (!location.ok()) {
+ return location.status();
+ }
+ // FlightServerBase::Init starts serving immediately; the blocking
Serve lifecycle wrapper
+ // is unnecessary for an in-process unit test.
+ return Init(arrow::flight::FlightServerOptions(*location));
+ }
+
+ ~BlockingFlightServer() override {
+ release();
+ static_cast<void>(Shutdown());
+ }
+
+ arrow::Status DoGet(const arrow::flight::ServerCallContext& context,
+ const arrow::flight::Ticket&,
+ std::unique_ptr<arrow::flight::FlightDataStream>*
stream) override {
+ if (_mode == Mode::DO_GET) {
+ {
+ std::lock_guard lock(_mutex);
+ _entered = true;
+ }
+ _cv.notify_all();
+ std::unique_lock lock(_mutex);
+ while (!_released && !context.is_cancelled()) {
+ _cv.wait_for(lock, std::chrono::milliseconds(5));
+ }
+ if (context.is_cancelled()) {
+ return arrow::Status::Cancelled("client cancelled blocked
DoGet");
+ }
+ }
+ *stream =
std::make_unique<arrow::flight::RecordBatchStream>(_batch_reader);
+ return arrow::Status::OK();
+ }
+
+ bool wait_until_entered(std::chrono::milliseconds timeout) {
+ if (_mode == Mode::NEXT) {
+ return _batch_reader->wait_until_entered(timeout);
+ }
+ std::unique_lock lock(_mutex);
+ return _cv.wait_for(lock, timeout, [this] { return _entered; });
+ }
+
+ void release() {
+ {
+ std::lock_guard lock(_mutex);
+ _released = true;
+ }
+ _cv.notify_all();
+ _batch_reader->release();
+ }
+
+private:
+ Mode _mode;
+ std::shared_ptr<BlockingRecordBatchReader> _batch_reader;
+ std::mutex _mutex;
+ std::condition_variable _cv;
+ bool _entered = false;
+ bool _released = false;
+};
+
+TFileRangeDesc remote_doris_range(const BlockingFlightServer& server) {
+ auto range = remote_doris_range();
+ auto& params = range.table_format_params.remote_doris_params;
+ params.__set_location_uri("grpc://localhost:" +
std::to_string(server.port()));
+ arrow::flight::Ticket ticket {.ticket = "ticket"};
+ params.__set_ticket(ticket.SerializeToString().ValueOrDie());
+ return range;
+}
+
std::vector<SlotDescriptor*> remote_slots(ObjectPool* pool, DescriptorTbl**
desc_tbl) {
DescriptorTblBuilder builder(pool);
builder.declare_tuple() <<
std::make_tuple(std::make_shared<DataTypeInt32>(), std::string("id"))
@@ -196,6 +319,16 @@ std::unique_ptr<RemoteDorisFileReader> create_reader(
std::move(factory));
}
+std::unique_ptr<RemoteDorisFileReader> create_flight_reader(
+ RuntimeProfile* profile, const TFileRangeDesc& range,
+ const std::vector<SlotDescriptor*>& slots,
std::shared_ptr<io::IOContext> io_ctx) {
+ auto system_properties = std::make_shared<io::FileSystemProperties>();
+ auto file_description = std::make_unique<io::FileDescription>();
+ file_description->path = "/dummyPath";
+ return std::make_unique<RemoteDorisFileReader>(system_properties,
file_description,
+ std::move(io_ctx), profile,
range, slots);
+}
+
Block make_request_block(const std::vector<ColumnDefinition>& schema,
const std::vector<int32_t>& local_ids) {
Block block;
@@ -467,4 +600,140 @@ TEST(RemoteDorisV2ReaderTest,
RejectsInvalidRemoteDorisRange) {
EXPECT_FALSE(reader->init(&state).ok());
}
+TEST(RemoteDorisV2ReaderTest, RuntimeCancellationInterruptsBlockedFlightDoGet)
{
+ BlockingFlightServer server(BlockingFlightServer::Mode::DO_GET);
+ const auto server_status = server.start();
+ ASSERT_TRUE(server_status.ok()) << server_status;
+ ObjectPool pool;
+ DescriptorTbl* desc_tbl = nullptr;
+ const auto slots = remote_slots(&pool, &desc_tbl);
+ RuntimeState state;
+ RuntimeProfile profile("remote_doris_v2_blocked_doget_test");
+ auto io_ctx = std::make_shared<io::IOContext>();
+ auto reader = create_flight_reader(&profile, remote_doris_range(server),
slots, io_ctx);
+ ASSERT_TRUE(reader->init(&state).ok());
+ auto request = std::make_shared<FileScanRequest>();
+ FileScanRequestBuilder builder(request.get());
+ ASSERT_TRUE(builder.add_non_predicate_column(LocalColumnId(0)).ok());
+
+ auto open_result =
+ std::async(std::launch::async, [&] { return
reader->open(std::move(request)); });
+ const bool entered = server.wait_until_entered(std::chrono::seconds(2));
+ if (!entered &&
+ open_result.wait_for(std::chrono::milliseconds(0)) ==
std::future_status::ready) {
+ FAIL() << "Flight DoGet failed before reaching the server: " <<
open_result.get();
+ }
+ ASSERT_TRUE(entered);
+ state.cancel(Status::Cancelled("cancel blocked Flight DoGet"));
+ const bool interrupted =
+ open_result.wait_for(std::chrono::milliseconds(750)) ==
std::future_status::ready;
+ if (!interrupted) {
+ server.release();
+ }
+ EXPECT_TRUE(interrupted);
+ EXPECT_FALSE(open_result.get().ok());
+}
+
+TEST(RemoteDorisV2ReaderTest,
BlockedDoGetRetainsQueryResourcesUntilWorkerExits) {
+ BlockingFlightServer server(BlockingFlightServer::Mode::DO_GET);
+ const auto server_status = server.start();
+ ASSERT_TRUE(server_status.ok()) << server_status;
+ ObjectPool pool;
+ DescriptorTbl* desc_tbl = nullptr;
+ const auto slots = remote_slots(&pool, &desc_tbl);
+ auto state = std::make_unique<MockRuntimeState>();
+ std::weak_ptr<ResourceContext> resource_ctx =
state->get_query_ctx()->resource_ctx();
+ RuntimeProfile profile("remote_doris_v2_doget_resource_context_test");
+ auto reader = create_flight_reader(&profile, remote_doris_range(server),
slots,
+ std::make_shared<io::IOContext>());
+ ASSERT_TRUE(reader->init(state.get()).ok());
+ auto request = std::make_shared<FileScanRequest>();
+ FileScanRequestBuilder builder(request.get());
+ ASSERT_TRUE(builder.add_non_predicate_column(LocalColumnId(0)).ok());
+
+ auto open_result =
+ std::async(std::launch::async, [&] { return
reader->open(std::move(request)); });
+ ASSERT_TRUE(server.wait_until_entered(std::chrono::seconds(2)));
+ state->cancel(Status::Cancelled("cancel blocked Flight DoGet"));
+ ASSERT_EQ(open_result.wait_for(std::chrono::milliseconds(750)),
std::future_status::ready);
+ EXPECT_FALSE(open_result.get().ok());
+
+ reader.reset();
+ state.reset();
+ // A detached DoGet still owns query-scoped Arrow objects, so it must
retain the matching
+ // resource context until those objects are released on the worker.
+ EXPECT_FALSE(resource_ctx.expired());
+ server.release();
+ for (int retries = 0; retries < 100 && !resource_ctx.expired(); ++retries)
{
+ std::this_thread::sleep_for(std::chrono::milliseconds(5));
+ }
+ EXPECT_TRUE(resource_ctx.expired());
+}
+
+TEST(RemoteDorisV2ReaderTest, ScannerStopInterruptsBlockedFlightNext) {
+ BlockingFlightServer server(BlockingFlightServer::Mode::NEXT);
+ const auto server_status = server.start();
+ ASSERT_TRUE(server_status.ok()) << server_status;
+ ObjectPool pool;
+ DescriptorTbl* desc_tbl = nullptr;
+ const auto slots = remote_slots(&pool, &desc_tbl);
+ RuntimeState state;
+ RuntimeProfile profile("remote_doris_v2_blocked_next_test");
+ auto io_ctx = std::make_shared<io::IOContext>();
+ auto reader = create_flight_reader(&profile, remote_doris_range(server),
slots, io_ctx);
+ ASSERT_TRUE(reader->init(&state).ok());
+ std::vector<ColumnDefinition> schema;
+ ASSERT_TRUE(reader->get_schema(&schema).ok());
+ auto request = std::make_shared<FileScanRequest>();
+ FileScanRequestBuilder builder(request.get());
+ ASSERT_TRUE(builder.add_non_predicate_column(LocalColumnId(0)).ok());
+ const auto open_status = reader->open(std::move(request));
+ ASSERT_TRUE(open_status.ok()) << open_status;
+
+ auto block = make_request_block(schema, {0});
+ size_t rows = 0;
+ bool eof = false;
+ auto next_result =
+ std::async(std::launch::async, [&] { return
reader->get_block(&block, &rows, &eof); });
+ ASSERT_TRUE(server.wait_until_entered(std::chrono::seconds(2)));
+ io_ctx->should_stop = true;
+ const bool interrupted =
+ next_result.wait_for(std::chrono::milliseconds(750)) ==
std::future_status::ready;
+ if (!interrupted) {
+ server.release();
+ }
+ EXPECT_TRUE(interrupted);
+ const auto next_status = next_result.get();
+ EXPECT_TRUE(!next_status.ok() || (rows == 0 && eof));
+ server.release();
+}
+
+TEST(RemoteDorisV2ReaderTest, CancellationStopsBeforeFlightNext) {
+ ObjectPool pool;
+ DescriptorTbl* desc_tbl = nullptr;
+ const auto slots = remote_slots(&pool, &desc_tbl);
+ RuntimeState state;
+ RuntimeProfile profile("remote_doris_v2_reader_cancel_test");
+ auto close_count = std::make_shared<int>(0);
+ auto io_ctx = std::make_shared<io::IOContext>();
+ auto reader = create_reader(&profile, remote_doris_range(), slots,
{make_batch({"id"})},
+ close_count, io_ctx);
+ ASSERT_TRUE(reader->init(&state).ok());
+ std::vector<ColumnDefinition> schema;
+ ASSERT_TRUE(reader->get_schema(&schema).ok());
+ auto request = std::make_shared<FileScanRequest>();
+ FileScanRequestBuilder builder(request.get());
+ ASSERT_TRUE(builder.add_non_predicate_column(LocalColumnId(0)).ok());
+ ASSERT_TRUE(reader->open(request).ok());
+
+ io_ctx->should_stop = true;
+ auto block = make_request_block(schema, {0});
+ size_t rows = 99;
+ bool eof = false;
+ ASSERT_TRUE(reader->get_block(&block, &rows, &eof).ok());
+ EXPECT_EQ(rows, 0);
+ EXPECT_TRUE(eof);
+ EXPECT_EQ(*close_count, 1);
+}
+
} // namespace doris::format::remote_doris
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]