From 6e94309a6cae89956f8a2b00b995dfffeaa3eeb4 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Wed, 8 Jul 2026 11:39:06 +0800 Subject: [PATCH 1/2] [fix](be) Release deletion vector rows on read errors ### What problem does this PR solve? Issue Number: None Related PR: #65351 Problem Summary: Deletion vector and position-delete cache builders allocated rows/maps with raw `new` and returned `nullptr` on read or parse failures. Those error paths leaked partially built cache values because `KVCache::get` only takes ownership when the builder returns a non-null pointer. Use RAII while building `DeleteRows` and `DeleteFile`, and release ownership only after successful construction. Add fault-injection/regression coverage for V1 Iceberg deletion-vector and position-delete failures, V1 Paimon Parquet/ORC deletion-vector failures, the shared Iceberg helper, and the format v2 TableReader path. ### Release note None ### Check List (For Author) - Test: Regression test / Unit Test - Added BE unit tests for V1/shared Iceberg deletion vector helper fault and stress paths. - Added BE unit tests for V1 Iceberg reader deletion vector and position-delete read failure cache-builder paths. - Added BE unit tests for V1 Paimon Parquet/ORC reader deletion vector read failure cache-builder paths. - Added BE unit tests for format v2 Iceberg deletion vector fault paths. - Extended Iceberg deletion vector regression to cover both scanner v1 and scanner v2 debug point fault injection. - Ran build-support/clang-format.sh locally. - Ran PATH=/usr/local/opt/llvm@16/bin:$PATH build-support/check-format.sh locally. - Ran git diff --check locally. - On gabriel@10.26.20.3 in /mnt/disk3/gabriel/Workspace/dev1/doris_pr65351_ut, ran ./run-be-ut.sh --run --filter="IcebergReaderTest.v1_position_delete_read_error_releases_cache_entry:IcebergReaderTest.v1_deletion_vector_read_error_releases_cache_entry:PaimonDeletionVectorTest.V1*:IcebergDeleteFileReaderHelperTest.*:IcebergV2ReaderTest.IcebergDeletionVector*:IcebergV2ReaderTest.IcebergTableReader*DeletionVector*"; 25 tests passed. - Attempted build-support/run-clang-tidy.sh --build-dir be/cmake-build-debug-dev-perf locally with CLANG_TIDY_BINARY=/usr/local/opt/llvm@16/bin/clang-tidy; blocked by the local macOS bash 3.2 missing mapfile. - Behavior changed: No - Does this need documentation: No --- .../iceberg_delete_file_reader_helper.cpp | 7 +- be/src/format/table/iceberg_reader_mixin.h | 26 +- be/src/format/table/paimon_reader.cpp | 13 +- be/src/format/table/paimon_reader.h | 4 + be/src/format_v2/table_reader.cpp | 16 +- ...iceberg_delete_file_reader_helper_test.cpp | 242 ++++++++++++++++++ .../table/iceberg/iceberg_reader_test.cpp | 68 +++++ .../format/table/paimon_cpp_reader_test.cpp | 73 ++++++ .../format_v2/table/iceberg_reader_test.cpp | 192 ++++++++++++++ .../test_iceberg_deletion_vector.groovy | 29 ++- 10 files changed, 652 insertions(+), 18 deletions(-) diff --git a/be/src/format/table/iceberg_delete_file_reader_helper.cpp b/be/src/format/table/iceberg_delete_file_reader_helper.cpp index 5845e440d95287..21bf9865868f13 100644 --- a/be/src/format/table/iceberg_delete_file_reader_helper.cpp +++ b/be/src/format/table/iceberg_delete_file_reader_helper.cpp @@ -44,6 +44,7 @@ #include "io/hdfs_builder.h" #include "runtime/runtime_state.h" #include "storage/predicate/column_predicate.h" +#include "util/debug_points.h" namespace doris { @@ -301,6 +302,10 @@ Status read_iceberg_deletion_vector(const TIcebergDeleteFileDesc& delete_file, if (!delete_file.__isset.content_offset || !delete_file.__isset.content_size_in_bytes) { return Status::InternalError("Deletion vector is missing content offset or length"); } + DBUG_EXECUTE_IF("IcebergDeleteFileReader.read_deletion_vector.io_error", + { return Status::IOError("injected Iceberg deletion vector read failure"); }); + DBUG_EXECUTE_IF("IcebergDeleteFileReader.read_deletion_vector.should_stop", + { return Status::EndOfFile("stop read."); }); TFileRangeDesc delete_range = build_iceberg_delete_file_range(delete_file.path); if (options.fs_name != nullptr && !options.fs_name->empty()) { @@ -316,7 +321,7 @@ Status read_iceberg_deletion_vector(const TIcebergDeleteFileDesc& delete_file, std::vector buf(delete_range.size); RETURN_IF_ERROR(dv_reader.read_at(delete_range.start_offset, {buf.data(), cast_set(delete_range.size)})); - return decode_iceberg_deletion_vector_buffer(buf.data(), delete_range.size, rows_to_delete); + return decode_deletion_vector_buffer(buf.data(), delete_range.size, rows_to_delete); } Status decode_iceberg_deletion_vector_buffer(const char* buf, size_t buffer_size, diff --git a/be/src/format/table/iceberg_reader_mixin.h b/be/src/format/table/iceberg_reader_mixin.h index f27fa83467c626..2e9bd06b6fdb48 100644 --- a/be/src/format/table/iceberg_reader_mixin.h +++ b/be/src/format/table/iceberg_reader_mixin.h @@ -19,6 +19,7 @@ #include #include +#include #include #include #include @@ -108,6 +109,16 @@ class IcebergReaderMixin : public BaseReader, public TableSchemaChangeHelper { _create_topn_row_id_column_iterator = create_func; } + Status TEST_read_deletion_vector(const std::string& data_file_path, + const TIcebergDeleteFileDesc& delete_file_desc) { + return read_deletion_vector(data_file_path, delete_file_desc); + } + + Status TEST_position_delete_base(const std::string& data_file_path, + const std::vector& delete_files) { + return _position_delete_base(data_file_path, delete_files); + } + protected: // ---- Hook implementations ---- @@ -627,18 +638,19 @@ Status IcebergReaderMixin::_position_delete_base( Status create_status = Status::OK(); auto* delete_file_cache = _kv_cache->template get( _delet_file_cache_key(delete_file.path), [&]() -> DeleteFile* { - auto* position_delete = new DeleteFile; + auto position_delete = std::make_unique(); TFileRangeDesc delete_file_range; delete_file_range.__set_fs_name(this->get_scan_range().fs_name); delete_file_range.path = delete_file.path; delete_file_range.start_offset = 0; delete_file_range.size = -1; delete_file_range.file_size = -1; - create_status = _read_position_delete_file(&delete_file_range, position_delete); + create_status = + _read_position_delete_file(&delete_file_range, position_delete.get()); if (!create_status) { return nullptr; } - return position_delete; + return position_delete.release(); }); if (create_status.is()) { continue; @@ -660,10 +672,10 @@ Status IcebergReaderMixin::_position_delete_base( SCOPED_TIMER(_iceberg_profile.delete_rows_sort_time); _iceberg_delete_rows = _kv_cache->template get(data_file_path, [&]() -> DeleteRows* { - auto* data_file_position_delete = new DeleteRows; + auto data_file_position_delete = std::make_unique(); _sort_delete_rows(delete_rows_array, num_delete_rows, *data_file_position_delete); - return data_file_position_delete; + return data_file_position_delete.release(); }); set_delete_rows(); COUNTER_UPDATE(_iceberg_profile.num_delete_rows, num_delete_rows); @@ -808,7 +820,7 @@ Status IcebergReaderMixin::read_deletion_vector( _iceberg_delete_rows = _kv_cache->template get( build_iceberg_deletion_vector_cache_key(data_file_path, delete_file_desc), [&]() -> DeleteRows* { - auto* delete_rows = new DeleteRows; + auto delete_rows = std::make_unique(); roaring::Roaring64Map bitmap; IcebergDeleteFileReaderOptions options; @@ -828,7 +840,7 @@ Status IcebergReaderMixin::read_deletion_vector( delete_rows->push_back(*it); } COUNTER_UPDATE(_iceberg_profile.num_delete_rows, delete_rows->size()); - return delete_rows; + return delete_rows.release(); }); RETURN_IF_ERROR(create_status); diff --git a/be/src/format/table/paimon_reader.cpp b/be/src/format/table/paimon_reader.cpp index 98f85658e32e0c..51a80f3a9675c8 100644 --- a/be/src/format/table/paimon_reader.cpp +++ b/be/src/format/table/paimon_reader.cpp @@ -20,6 +20,7 @@ #include #include +#include #include #include "common/status.h" @@ -134,7 +135,7 @@ Status PaimonOrcReader::_init_deletion_vector() { using DeleteRows = std::vector; _delete_rows = _kv_cache->get( build_paimon_deletion_vector_cache_key(deletion_file), [&]() -> DeleteRows* { - auto* delete_rows = new DeleteRows; + auto delete_rows = std::make_unique(); TFileRangeDesc delete_range; delete_range.__set_fs_name(get_scan_range().fs_name); @@ -160,12 +161,12 @@ Status PaimonOrcReader::_init_deletion_vector() { SCOPED_TIMER(_paimon_profile.parse_deletion_vector_time); create_status = decode_paimon_deletion_vector_buffer(buffer.data(), bytes_read, - delete_rows); + delete_rows.get()); if (!create_status.ok()) [[unlikely]] { return nullptr; } COUNTER_UPDATE(_paimon_profile.num_delete_rows, delete_rows->size()); - return delete_rows; + return delete_rows.release(); }); RETURN_IF_ERROR(create_status); if (!_delete_rows->empty()) [[likely]] { @@ -233,7 +234,7 @@ Status PaimonParquetReader::_init_deletion_vector() { using DeleteRows = std::vector; _delete_rows = _kv_cache->get( build_paimon_deletion_vector_cache_key(deletion_file), [&]() -> DeleteRows* { - auto* delete_rows = new DeleteRows; + auto delete_rows = std::make_unique(); TFileRangeDesc delete_range; delete_range.__set_fs_name(get_scan_range().fs_name); @@ -259,12 +260,12 @@ Status PaimonParquetReader::_init_deletion_vector() { SCOPED_TIMER(_paimon_profile.parse_deletion_vector_time); create_status = decode_paimon_deletion_vector_buffer(buffer.data(), bytes_read, - delete_rows); + delete_rows.get()); if (!create_status.ok()) [[unlikely]] { return nullptr; } COUNTER_UPDATE(_paimon_profile.num_delete_rows, delete_rows->size()); - return delete_rows; + return delete_rows.release(); }); RETURN_IF_ERROR(create_status); if (!_delete_rows->empty()) [[likely]] { diff --git a/be/src/format/table/paimon_reader.h b/be/src/format/table/paimon_reader.h index 10da0b784e2b46..616c58275b30e2 100644 --- a/be/src/format/table/paimon_reader.h +++ b/be/src/format/table/paimon_reader.h @@ -64,6 +64,8 @@ class PaimonOrcReader final : public OrcReader, public TableSchemaChangeHelper { } ~PaimonOrcReader() final = default; + Status TEST_init_deletion_vector() { return _init_deletion_vector(); } + protected: Status on_before_init_reader(ReaderInitContext* ctx) override; @@ -109,6 +111,8 @@ class PaimonParquetReader final : public ParquetReader, public TableSchemaChange } ~PaimonParquetReader() final = default; + Status TEST_init_deletion_vector() { return _init_deletion_vector(); } + protected: Status on_before_init_reader(ReaderInitContext* ctx) override; diff --git a/be/src/format_v2/table_reader.cpp b/be/src/format_v2/table_reader.cpp index e1bd6b04ad2fff..99cacbc7a68ee0 100644 --- a/be/src/format_v2/table_reader.cpp +++ b/be/src/format_v2/table_reader.cpp @@ -22,6 +22,7 @@ #include #include +#include #include #include #include @@ -46,6 +47,7 @@ #include "format_v2/native/native_reader.h" #include "format_v2/parquet/parquet_reader.h" #include "storage/segment/condition_cache.h" +#include "util/debug_points.h" #include "util/string_util.h" namespace doris::format { @@ -765,7 +767,7 @@ Status TableReader::_parse_delete_predicates(const SplitReadOptions& options) { Status create_status = Status::OK(); _delete_rows = options.cache->get(desc.key, [&]() -> DeleteRows* { - auto* delete_rows = new DeleteRows; + auto delete_rows = std::make_unique(); DeletionVectorReader dv_reader(_runtime_state, _scanner_profile, *_scan_params, desc, _io_ctx.get()); @@ -776,6 +778,14 @@ Status TableReader::_parse_delete_predicates(const SplitReadOptions& options) { size_t bytes_read = desc.size; std::vector buffer(bytes_read); + DBUG_EXECUTE_IF("TableReader.parse_deletion_vector.io_error", { + create_status = Status::IOError("injected format v2 deletion vector read failure"); + return nullptr; + }); + DBUG_EXECUTE_IF("TableReader.parse_deletion_vector.should_stop", { + create_status = Status::EndOfFile("stop read."); + return nullptr; + }); create_status = dv_reader.read_at(desc.start_offset, {buffer.data(), bytes_read}); if (!create_status.ok()) [[unlikely]] { return nullptr; @@ -783,12 +793,12 @@ Status TableReader::_parse_delete_predicates(const SplitReadOptions& options) { const char* buf = buffer.data(); SCOPED_TIMER(_profile.parse_delete_file_time); - create_status = parse_deletion_vector(buf, bytes_read, desc.format, delete_rows); + create_status = parse_deletion_vector(buf, bytes_read, desc.format, delete_rows.get()); if (!create_status.ok()) [[unlikely]] { return nullptr; } COUNTER_UPDATE(_profile.num_delete_rows, delete_rows->size()); - return delete_rows; + return delete_rows.release(); }); RETURN_IF_ERROR(create_status); } diff --git a/be/test/format/table/iceberg/iceberg_delete_file_reader_helper_test.cpp b/be/test/format/table/iceberg/iceberg_delete_file_reader_helper_test.cpp index 6ce562d6de0d06..34d2e60f712a5a 100644 --- a/be/test/format/table/iceberg/iceberg_delete_file_reader_helper_test.cpp +++ b/be/test/format/table/iceberg/iceberg_delete_file_reader_helper_test.cpp @@ -17,15 +17,24 @@ #include "format/table/iceberg_delete_file_reader_helper.h" +#include #include #include +#include +#include +#include +#include +#include #include #include +#include "exec/common/endian.h" #include "io/fs/file_meta_cache.h" +#include "roaring/roaring64map.hh" #include "runtime/runtime_profile.h" #include "runtime/runtime_state.h" +#include "runtime/thread_context.h" namespace doris { @@ -49,6 +58,58 @@ class CollectPositionDeleteVisitor final : public IcebergPositionDeleteVisitor { size_t total_rows = 0; }; +TFileScanRangeParams make_local_parquet_scan_params() { + TFileScanRangeParams scan_params; + scan_params.__set_file_type(TFileType::FILE_LOCAL); + scan_params.__set_format_type(TFileFormatType::FORMAT_PARQUET); + return scan_params; +} + +TIcebergDeleteFileDesc make_iceberg_deletion_vector(const std::string& path, int64_t offset, + int64_t size) { + TIcebergDeleteFileDesc delete_file; + delete_file.__set_content(3); + delete_file.__set_path(path); + delete_file.__set_content_offset(offset); + delete_file.__set_content_size_in_bytes(size); + return delete_file; +} + +int64_t write_iceberg_deletion_vector_file(const std::string& file_path, + const std::vector& deleted_positions) { + roaring::Roaring64Map rows; + for (const auto position : deleted_positions) { + rows.add(position); + } + + const size_t bitmap_size = rows.getSizeInBytes(); + std::vector blob(4 + 4 + bitmap_size + 4); + rows.write(blob.data() + 8); + + const uint32_t total_length = static_cast(4 + bitmap_size); + BigEndian::Store32(blob.data(), total_length); + constexpr char DV_MAGIC[] = {'\xD1', '\xD3', '\x39', '\x64'}; + memcpy(blob.data() + 4, DV_MAGIC, 4); + BigEndian::Store32(blob.data() + 8 + bitmap_size, 0); + + std::ofstream output(file_path, std::ios::binary); + EXPECT_TRUE(output.is_open()); + output.write(blob.data(), static_cast(blob.size())); + EXPECT_TRUE(output.good()); + return static_cast(blob.size()); +} + +IcebergDeleteFileReaderOptions make_delete_file_reader_options( + RuntimeState* state, RuntimeProfile* profile, const TFileScanRangeParams* scan_params, + io::IOContext* io_ctx) { + return { + .state = state, + .profile = profile, + .scan_params = scan_params, + .io_ctx = io_ctx, + }; +} + } // namespace TEST(IcebergDeleteFileReaderHelperTest, BuildDeleteFileRange) { @@ -119,6 +180,187 @@ TEST(IcebergDeleteFileReaderHelperTest, DeletionVectorCacheKeyEscapesPathBoundar build_iceberg_deletion_vector_cache_key(second_data_file_path, second_delete_file)); } +TEST(IcebergDeleteFileReaderHelperTest, ReadDeletionVectorReportsMissingFile) { + const auto test_dir = + std::filesystem::temp_directory_path() / "doris_iceberg_deletion_vector_missing_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + + const auto missing_path = (test_dir / "missing-delete-vector.bin").string(); + + RuntimeProfile profile("test_profile"); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + auto scan_params = make_local_parquet_scan_params(); + IcebergDeleteFileIOContext io_context(&state); + roaring::Roaring64Map rows_to_delete; + auto options = + make_delete_file_reader_options(&state, &profile, &scan_params, &io_context.io_ctx); + + auto status = read_iceberg_deletion_vector(make_iceberg_deletion_vector(missing_path, 0, 16), + options, &rows_to_delete); + + EXPECT_FALSE(status.ok()); + EXPECT_NE(status.to_string().find(missing_path), std::string::npos); + EXPECT_EQ(rows_to_delete.cardinality(), 0); + std::filesystem::remove_all(test_dir); +} + +TEST(IcebergDeleteFileReaderHelperTest, ReadDeletionVectorReportsShortRead) { + const auto test_dir = std::filesystem::temp_directory_path() / + "doris_iceberg_deletion_vector_short_read_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + + const auto dv_path = (test_dir / "delete-vector.bin").string(); + const auto dv_size = write_iceberg_deletion_vector_file(dv_path, {1, 3}); + + RuntimeProfile profile("test_profile"); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + auto scan_params = make_local_parquet_scan_params(); + IcebergDeleteFileIOContext io_context(&state); + roaring::Roaring64Map rows_to_delete; + auto options = + make_delete_file_reader_options(&state, &profile, &scan_params, &io_context.io_ctx); + + auto status = read_iceberg_deletion_vector( + make_iceberg_deletion_vector(dv_path, 0, dv_size + 1), options, &rows_to_delete); + + EXPECT_FALSE(status.ok()); + EXPECT_EQ(rows_to_delete.cardinality(), 0); + std::filesystem::remove_all(test_dir); +} + +TEST(IcebergDeleteFileReaderHelperTest, ReadDeletionVectorStopsWhenIoContextStops) { + const auto test_dir = + std::filesystem::temp_directory_path() / "doris_iceberg_deletion_vector_stop_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + + const auto dv_path = (test_dir / "delete-vector.bin").string(); + const auto dv_size = write_iceberg_deletion_vector_file(dv_path, {0}); + + RuntimeProfile profile("test_profile"); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + auto scan_params = make_local_parquet_scan_params(); + IcebergDeleteFileIOContext io_context(&state); + io_context.io_ctx.should_stop = true; + roaring::Roaring64Map rows_to_delete; + auto options = + make_delete_file_reader_options(&state, &profile, &scan_params, &io_context.io_ctx); + + auto status = read_iceberg_deletion_vector(make_iceberg_deletion_vector(dv_path, 0, dv_size), + options, &rows_to_delete); + + EXPECT_TRUE(status.is()); + EXPECT_NE(status.to_string().find("stop read"), std::string::npos); + EXPECT_EQ(rows_to_delete.cardinality(), 0); + std::filesystem::remove_all(test_dir); +} + +TEST(IcebergDeleteFileReaderHelperTest, DecodeDeletionVectorRejectsCorruptPayload) { + std::vector corrupted(12, 0); + BigEndian::Store32(corrupted.data(), 4); + memcpy(corrupted.data() + 4, "BAD!", 4); + + roaring::Roaring64Map rows_to_delete; + auto status = decode_iceberg_deletion_vector_buffer(corrupted.data(), corrupted.size(), + &rows_to_delete); + + EXPECT_TRUE(status.is()); + EXPECT_NE(status.to_string().find("magic number mismatch"), std::string::npos); + EXPECT_EQ(rows_to_delete.cardinality(), 0); +} + +TEST(IcebergDeleteFileReaderHelperTest, ReadDeletionVectorReadsMillionDeletePositions) { + const auto test_dir = + std::filesystem::temp_directory_path() / "doris_iceberg_deletion_vector_large_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + + constexpr uint64_t delete_position_count = 1'000'000; + std::vector deleted_positions; + deleted_positions.reserve(delete_position_count); + for (uint64_t position = 0; position < delete_position_count; ++position) { + deleted_positions.push_back(position * 2); + } + + const auto dv_path = (test_dir / "delete-vector.bin").string(); + const auto dv_size = write_iceberg_deletion_vector_file(dv_path, deleted_positions); + + RuntimeProfile profile("test_profile"); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + auto scan_params = make_local_parquet_scan_params(); + IcebergDeleteFileIOContext io_context(&state); + roaring::Roaring64Map rows_to_delete; + auto options = + make_delete_file_reader_options(&state, &profile, &scan_params, &io_context.io_ctx); + + ASSERT_TRUE(read_iceberg_deletion_vector(make_iceberg_deletion_vector(dv_path, 0, dv_size), + options, &rows_to_delete) + .ok()); + + EXPECT_EQ(rows_to_delete.cardinality(), delete_position_count); + EXPECT_TRUE(rows_to_delete.contains(static_cast(0))); + EXPECT_TRUE(rows_to_delete.contains(static_cast(999'998))); + EXPECT_TRUE(rows_to_delete.contains(static_cast(1'999'998))); + EXPECT_FALSE(rows_to_delete.contains(static_cast(1'999'999))); + std::filesystem::remove_all(test_dir); +} + +TEST(IcebergDeleteFileReaderHelperTest, ReadDeletionVectorConcurrentReadsAreStable) { + const auto test_dir = std::filesystem::temp_directory_path() / + "doris_iceberg_deletion_vector_concurrent_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + + const auto first_dv_path = (test_dir / "delete-vector-first.bin").string(); + const auto second_dv_path = (test_dir / "delete-vector-second.bin").string(); + const auto first_dv_size = write_iceberg_deletion_vector_file(first_dv_path, {0, 2, 4}); + const auto second_dv_size = write_iceberg_deletion_vector_file(second_dv_path, {1, 3, 5}); + + auto read_and_verify = [](const std::string& path, int64_t size, + const std::vector& expected_positions) -> Status { + RuntimeProfile profile("test_profile"); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + auto scan_params = make_local_parquet_scan_params(); + IcebergDeleteFileIOContext io_context(&state); + roaring::Roaring64Map rows_to_delete; + auto options = + make_delete_file_reader_options(&state, &profile, &scan_params, &io_context.io_ctx); + RETURN_IF_ERROR(read_iceberg_deletion_vector(make_iceberg_deletion_vector(path, 0, size), + options, &rows_to_delete)); + if (rows_to_delete.cardinality() != expected_positions.size()) { + return Status::InternalError("unexpected deletion vector cardinality"); + } + for (const auto position : expected_positions) { + if (!rows_to_delete.contains(position)) { + return Status::InternalError("missing deletion vector position {}", position); + } + } + return Status::OK(); + }; + + std::vector workers; + std::vector statuses(16); + for (size_t idx = 0; idx < statuses.size(); ++idx) { + workers.emplace_back([&, idx]() { + SCOPED_INIT_THREAD_CONTEXT(); + if (idx % 2 == 0) { + statuses[idx] = read_and_verify(first_dv_path, first_dv_size, {0, 2, 4}); + } else { + statuses[idx] = read_and_verify(second_dv_path, second_dv_size, {1, 3, 5}); + } + }); + } + for (auto& worker : workers) { + worker.join(); + } + for (const auto& status : statuses) { + EXPECT_TRUE(status.ok()) << status; + } + std::filesystem::remove_all(test_dir); +} + TEST(IcebergDeleteFileReaderHelperTest, ReadMixedEncodingParquetPositionDeleteFile) { RuntimeProfile profile("test_profile"); RuntimeState runtime_state((TQueryOptions()), TQueryGlobals()); diff --git a/be/test/format/table/iceberg/iceberg_reader_test.cpp b/be/test/format/table/iceberg/iceberg_reader_test.cpp index c261859731256a..89bf367d06b053 100644 --- a/be/test/format/table/iceberg/iceberg_reader_test.cpp +++ b/be/test/format/table/iceberg/iceberg_reader_test.cpp @@ -46,6 +46,7 @@ #include "core/data_type/data_type_string.h" #include "core/data_type/data_type_struct.h" #include "format/column_descriptor.h" +#include "format/format_common.h" #include "format/orc/vorc_reader.h" #include "format/parquet/vparquet_column_chunk_reader.h" #include "format/parquet/vparquet_reader.h" @@ -53,6 +54,7 @@ #include "io/fs/file_reader_writer_fwd.h" #include "io/fs/file_system.h" #include "io/fs/local_file_system.h" +#include "io/io_common.h" #include "runtime/descriptors.h" #include "runtime/runtime_state.h" #include "storage/olap_scan_common.h" @@ -777,6 +779,72 @@ TEST_F(IcebergReaderTest, read_iceberg_parquet_file) { verify_test_results(block, read_rows); } +TEST_F(IcebergReaderTest, v1_deletion_vector_read_error_releases_cache_entry) { + RuntimeState runtime_state = RuntimeState(TQueryOptions(), TQueryGlobals()); + TFileScanRangeParams scan_params; + scan_params.__set_file_type(TFileType::FILE_LOCAL); + scan_params.__set_format_type(TFileFormatType::FORMAT_PARQUET); + + TFileRangeDesc scan_range; + scan_range.__set_fs_name(""); + scan_range.__set_path("data.parquet"); + scan_range.__set_start_offset(0); + scan_range.__set_size(0); + + RuntimeProfile profile("test_profile"); + cctz::time_zone ctz; + TimezoneUtils::find_cctz_time_zone(TimezoneUtils::default_time_zone, ctz); + io::IOContext io_ctx; + ShardedKVCache kv_cache(8); + + IcebergParquetReader iceberg_reader(&kv_cache, &profile, scan_params, scan_range, 1024, &ctz, + &io_ctx, &runtime_state, cache.get()); + + TIcebergDeleteFileDesc delete_file; + delete_file.__set_content(IcebergReaderMixin::DELETION_VECTOR); + delete_file.__set_path("./be/test/exec/test_data/missing_iceberg_v1_delete_vector.bin"); + delete_file.__set_content_offset(0); + delete_file.__set_content_size_in_bytes(16); + + const auto status = + iceberg_reader.TEST_read_deletion_vector("file:///tmp/data.parquet", delete_file); + + ASSERT_FALSE(status.ok()); + EXPECT_NE(status.to_string().find(delete_file.path), std::string::npos); +} + +TEST_F(IcebergReaderTest, v1_position_delete_read_error_releases_cache_entry) { + RuntimeState runtime_state = RuntimeState(TQueryOptions(), TQueryGlobals()); + TFileScanRangeParams scan_params; + scan_params.__set_file_type(TFileType::FILE_LOCAL); + scan_params.__set_format_type(TFileFormatType::FORMAT_PARQUET); + + TFileRangeDesc scan_range; + scan_range.__set_fs_name(""); + scan_range.__set_path("data.parquet"); + scan_range.__set_start_offset(0); + scan_range.__set_size(0); + + RuntimeProfile profile("test_profile"); + cctz::time_zone ctz; + TimezoneUtils::find_cctz_time_zone(TimezoneUtils::default_time_zone, ctz); + io::IOContext io_ctx; + ShardedKVCache kv_cache(8); + + IcebergParquetReader iceberg_reader(&kv_cache, &profile, scan_params, scan_range, 1024, &ctz, + &io_ctx, &runtime_state, cache.get()); + + TIcebergDeleteFileDesc delete_file; + delete_file.__set_content(IcebergReaderMixin::POSITION_DELETE); + delete_file.__set_path("./be/test/exec/test_data/missing_iceberg_v1_position_delete.parquet"); + + const auto status = + iceberg_reader.TEST_position_delete_base("file:///tmp/data.parquet", {delete_file}); + + ASSERT_FALSE(status.ok()); + EXPECT_NE(status.to_string().find(delete_file.path), std::string::npos); +} + // Test reading real Iceberg Orc file using IcebergTableReader TEST_F(IcebergReaderTest, read_iceberg_orc_file) { // Read only: name, profile.address.coordinates.lat, profile.address.coordinates.lng, profile.contact.email diff --git a/be/test/format/table/paimon_cpp_reader_test.cpp b/be/test/format/table/paimon_cpp_reader_test.cpp index ecec461d6b120f..2858956c8f059b 100644 --- a/be/test/format/table/paimon_cpp_reader_test.cpp +++ b/be/test/format/table/paimon_cpp_reader_test.cpp @@ -17,6 +17,7 @@ #include "format/table/paimon_cpp_reader.h" +#include #include #include @@ -25,10 +26,14 @@ #include "core/block/block.h" #include "exec/common/endian.h" +#include "format/format_common.h" #include "format/table/paimon_reader.h" +#include "io/fs/file_meta_cache.h" +#include "io/io_common.h" #include "roaring/roaring.hh" #include "runtime/runtime_profile.h" #include "runtime/runtime_state.h" +#include "util/timezone_utils.h" namespace doris { @@ -50,6 +55,33 @@ std::vector build_paimon_deletion_vector_buffer(const std::vector #include #include +#include #include +#include "common/config.h" #include "common/consts.h" #include "core/assert_cast.h" #include "core/block/block.h" @@ -65,6 +67,7 @@ #include "runtime/runtime_profile.h" #include "runtime/runtime_state.h" #include "storage/segment/condition_cache.h" +#include "util/debug_points.h" namespace doris::format { namespace { @@ -426,6 +429,24 @@ int64_t write_iceberg_deletion_vector_file(const std::string& file_path, return static_cast(blob.size()); } +class ScopedDebugPoint { +public: + explicit ScopedDebugPoint(std::string name) + : _name(std::move(name)), _enable_debug_points(config::enable_debug_points) { + config::enable_debug_points = true; + DebugPoints::instance()->add(_name); + } + + ~ScopedDebugPoint() { + DebugPoints::instance()->remove(_name); + config::enable_debug_points = _enable_debug_points; + } + +private: + std::string _name; + bool _enable_debug_points; +}; + Block build_table_block(const std::vector& columns) { Block block; for (const auto& column : columns) { @@ -1258,6 +1279,28 @@ TEST(IcebergV2ReaderTest, IcebergDeletionVectorRejectsMultipleDeleteFiles) { EXPECT_FALSE(status.ok()); } +TEST(IcebergV2ReaderTest, IcebergDeletionVectorRejectsMissingRange) { + TIcebergDeleteFileDesc delete_file; + delete_file.__set_content(3); + delete_file.__set_path("dv.bin"); + + TTableFormatFileDesc table_format_desc; + TIcebergFileDesc iceberg_desc; + iceberg_desc.__set_format_version(2); + iceberg_desc.__set_delete_files({delete_file}); + table_format_desc.__set_iceberg_params(iceberg_desc); + + IcebergTableReaderDeleteFileTestHelper reader; + DeleteFileDesc desc; + bool has_delete_file = false; + auto status = reader.parse_deletion_vector_file(table_format_desc, &desc, &has_delete_file); + + EXPECT_FALSE(status.ok()); + EXPECT_TRUE(status.is()); + EXPECT_NE(status.to_string().find("missing content offset or length"), std::string::npos); + EXPECT_FALSE(has_delete_file); +} + TEST(IcebergV2ReaderTest, IcebergTableReaderAppliesDeletionVectorFile) { const auto test_dir = std::filesystem::temp_directory_path() / "doris_iceberg_deletion_vector_file_test"; @@ -1305,6 +1348,155 @@ TEST(IcebergV2ReaderTest, IcebergTableReaderAppliesDeletionVectorFile) { std::filesystem::remove_all(test_dir); } +TEST(IcebergV2ReaderTest, IcebergTableReaderReportsInjectedDeletionVectorReadError) { + const auto test_dir = std::filesystem::temp_directory_path() / + "doris_iceberg_v2_deletion_vector_injected_read_error_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + + const auto file_path = (test_dir / "split.parquet").string(); + const auto dv_path = (test_dir / "delete-vector.bin").string(); + write_int_pair_parquet_file(file_path, {1, 2, 3}, {10, 20, 30}, {"one", "two", "three"}); + const auto dv_size = write_iceberg_deletion_vector_file(dv_path, {0}); + + std::vector projected_columns; + projected_columns.push_back(make_table_column(0, "id", std::make_shared())); + + RuntimeProfile profile("test_profile"); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + auto scan_params = make_local_parquet_scan_params(); + io::FileReaderStats file_reader_stats; + io::FileCacheStatistics file_cache_stats; + auto io_ctx = make_io_context(&file_reader_stats, &file_cache_stats); + ShardedKVCache cache(1); + doris::format::iceberg::IcebergTableReader reader; + ASSERT_TRUE(reader.init({ + .projected_columns = projected_columns, + .conjuncts = {}, + .format = FileFormat::PARQUET, + .scan_params = &scan_params, + .io_ctx = io_ctx, + .runtime_state = &state, + .scanner_profile = &profile, + }) + .ok()); + + auto split_options = build_split_options(file_path); + split_options.cache = &cache; + split_options.current_range.__set_table_format_params(make_iceberg_table_format_desc( + file_path, {make_iceberg_deletion_vector(dv_path, 0, dv_size)})); + + ScopedDebugPoint debug_point("TableReader.parse_deletion_vector.io_error"); + auto status = reader.prepare_split(split_options); + + EXPECT_FALSE(status.ok()); + EXPECT_NE(status.to_string().find("injected format v2 deletion vector read failure"), + std::string::npos); + ASSERT_TRUE(reader.close().ok()); + std::filesystem::remove_all(test_dir); +} + +TEST(IcebergV2ReaderTest, IcebergTableReaderStopsDuringDeletionVectorRead) { + const auto test_dir = std::filesystem::temp_directory_path() / + "doris_iceberg_v2_deletion_vector_should_stop_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + + const auto file_path = (test_dir / "split.parquet").string(); + const auto dv_path = (test_dir / "delete-vector.bin").string(); + write_int_pair_parquet_file(file_path, {1, 2, 3}, {10, 20, 30}, {"one", "two", "three"}); + const auto dv_size = write_iceberg_deletion_vector_file(dv_path, {0}); + + std::vector projected_columns; + projected_columns.push_back(make_table_column(0, "id", std::make_shared())); + + RuntimeProfile profile("test_profile"); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + auto scan_params = make_local_parquet_scan_params(); + io::FileReaderStats file_reader_stats; + io::FileCacheStatistics file_cache_stats; + auto io_ctx = make_io_context(&file_reader_stats, &file_cache_stats); + ShardedKVCache cache(1); + doris::format::iceberg::IcebergTableReader reader; + ASSERT_TRUE(reader.init({ + .projected_columns = projected_columns, + .conjuncts = {}, + .format = FileFormat::PARQUET, + .scan_params = &scan_params, + .io_ctx = io_ctx, + .runtime_state = &state, + .scanner_profile = &profile, + }) + .ok()); + + auto split_options = build_split_options(file_path); + split_options.cache = &cache; + split_options.current_range.__set_table_format_params(make_iceberg_table_format_desc( + file_path, {make_iceberg_deletion_vector(dv_path, 0, dv_size)})); + + ScopedDebugPoint debug_point("TableReader.parse_deletion_vector.should_stop"); + auto status = reader.prepare_split(split_options); + + EXPECT_TRUE(status.is()); + EXPECT_NE(status.to_string().find("stop read"), std::string::npos); + ASSERT_TRUE(reader.close().ok()); + std::filesystem::remove_all(test_dir); +} + +TEST(IcebergV2ReaderTest, IcebergTableReaderRejectsCorruptDeletionVectorPayload) { + const auto test_dir = + std::filesystem::temp_directory_path() / "doris_iceberg_v2_corrupt_dv_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + + const auto file_path = (test_dir / "split.parquet").string(); + const auto dv_path = (test_dir / "delete-vector.bin").string(); + write_int_pair_parquet_file(file_path, {1, 2, 3}, {10, 20, 30}, {"one", "two", "three"}); + std::vector corrupted(12, 0); + BigEndian::Store32(corrupted.data(), 4); + memcpy(corrupted.data() + 4, "BAD!", 4); + { + std::ofstream output(dv_path, std::ios::binary); + ASSERT_TRUE(output.is_open()); + output.write(corrupted.data(), static_cast(corrupted.size())); + ASSERT_TRUE(output.good()); + } + + std::vector projected_columns; + projected_columns.push_back(make_table_column(0, "id", std::make_shared())); + + RuntimeProfile profile("test_profile"); + RuntimeState state {TQueryOptions(), TQueryGlobals()}; + auto scan_params = make_local_parquet_scan_params(); + io::FileReaderStats file_reader_stats; + io::FileCacheStatistics file_cache_stats; + auto io_ctx = make_io_context(&file_reader_stats, &file_cache_stats); + ShardedKVCache cache(1); + doris::format::iceberg::IcebergTableReader reader; + ASSERT_TRUE(reader.init({ + .projected_columns = projected_columns, + .conjuncts = {}, + .format = FileFormat::PARQUET, + .scan_params = &scan_params, + .io_ctx = io_ctx, + .runtime_state = &state, + .scanner_profile = &profile, + }) + .ok()); + + auto split_options = build_split_options(file_path); + split_options.cache = &cache; + split_options.current_range.__set_table_format_params(make_iceberg_table_format_desc( + file_path, + {make_iceberg_deletion_vector(dv_path, 0, static_cast(corrupted.size()))})); + auto status = reader.prepare_split(split_options); + + EXPECT_TRUE(status.is()); + EXPECT_NE(status.to_string().find("magic number mismatch"), std::string::npos); + ASSERT_TRUE(reader.close().ok()); + std::filesystem::remove_all(test_dir); +} + TEST(IcebergV2ReaderTest, IcebergTableReaderDoesNotPushDownAggregateWithDeletes) { const auto test_dir = std::filesystem::temp_directory_path() / "doris_iceberg_aggregate_delete_test"; diff --git a/regression-test/suites/external_table_p0/iceberg/test_iceberg_deletion_vector.groovy b/regression-test/suites/external_table_p0/iceberg/test_iceberg_deletion_vector.groovy index 63fcf1cfda0809..58abc32a34d65c 100644 --- a/regression-test/suites/external_table_p0/iceberg/test_iceberg_deletion_vector.groovy +++ b/regression-test/suites/external_table_p0/iceberg/test_iceberg_deletion_vector.groovy @@ -15,7 +15,7 @@ // specific language governing permissions and limitations // under the License. -suite("test_iceberg_deletion_vector", "p0,external") { +suite("test_iceberg_deletion_vector", "p0,external,nonConcurrent") { String enabled = context.config.otherConfigs.get("enableIcebergTest") if (enabled == null || !enabled.equalsIgnoreCase("true")) { logger.info("disable iceberg test.") @@ -328,6 +328,33 @@ class IcebergRestCatalog { qt_q3 """ SELECT * FROM dv_test_orc ORDER BY id; """ qt_no_delete """ SELECT * FROM dv_test_no_delete ORDER BY id; """ + def enableFileScannerV2Rows = sql """SHOW VARIABLES LIKE 'enable_file_scanner_v2'""" + assertTrue(enableFileScannerV2Rows.size() > 0, + "Session variable enable_file_scanner_v2 is not found") + String originalEnableFileScannerV2 = enableFileScannerV2Rows[0][1].toString() + try { + sql """set enable_file_scanner_v2=false""" + GetDebugPoint().clearDebugPointsForAllBEs() + GetDebugPoint().enableDebugPointForAllBEs( + "IcebergDeleteFileReader.read_deletion_vector.io_error") + test { + sql """ SELECT count(*) FROM dv_test; """ + exception "injected Iceberg deletion vector read failure" + } + + sql """set enable_file_scanner_v2=true""" + GetDebugPoint().clearDebugPointsForAllBEs() + GetDebugPoint().enableDebugPointForAllBEs( + "TableReader.parse_deletion_vector.io_error") + test { + sql """ SELECT count(*) FROM dv_test; """ + exception "injected format v2 deletion vector read failure" + } + } finally { + GetDebugPoint().clearDebugPointsForAllBEs() + sql """set enable_file_scanner_v2=${originalEnableFileScannerV2}""" + } + // Delete-type matrix checks cover equality-only, position-only, DV-only, DV+position, // and DV+equality combinations. qt_equality_only """ SELECT * FROM dv_delete_matrix_equality_only ORDER BY id; """ From 57f3721198c5c0d589ff4b8ad85fc166db5d496a Mon Sep 17 00:00:00 2001 From: lidongyang Date: Thu, 9 Jul 2026 13:34:17 +0800 Subject: [PATCH 2/2] reduce BE/FE memory limits for external regression pipeline - be.conf: add mem_limit=35% (was default 90%), raise max_sys_mem_available_low_water_mark_bytes from 66MB to 2GB - fe.conf: reduce -Xmx from 8192m to 4096m (actual heap usage at idle is ~2.1GB, so 4GB is ample for regression tests) Co-Authored-By: Claude --- regression-test/pipeline/external/conf/be.conf | 3 ++- regression-test/pipeline/external/conf/fe.conf | 2 +- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/regression-test/pipeline/external/conf/be.conf b/regression-test/pipeline/external/conf/be.conf index f29f342366e273..beea7be7175c38 100644 --- a/regression-test/pipeline/external/conf/be.conf +++ b/regression-test/pipeline/external/conf/be.conf @@ -53,7 +53,8 @@ priority_networks=172.19.0.0/24 enable_fuzzy_mode=true max_depth_of_expr_tree=200 enable_feature_binlog=true -max_sys_mem_available_low_water_mark_bytes=69206016 +mem_limit=35% +max_sys_mem_available_low_water_mark_bytes=2147483648 user_files_secure_path=/ enable_debug_points=true # debug scanner context dead loop diff --git a/regression-test/pipeline/external/conf/fe.conf b/regression-test/pipeline/external/conf/fe.conf index 365c0b9337576e..fc818a862a1a15 100644 --- a/regression-test/pipeline/external/conf/fe.conf +++ b/regression-test/pipeline/external/conf/fe.conf @@ -28,7 +28,7 @@ DATE = `date +%Y%m%d-%H%M%S` JAVA_OPTS="-Xmx4096m -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=$DORIS_HOME/log/fe.jmap -XX:+UseMembar -XX:SurvivorRatio=8 -XX:MaxTenuringThreshold=7 -XX:+PrintGCDateStamps -XX:+PrintGCDetails -XX:+PrintClassHistogramAfterFullGC -XX:+UseConcMarkSweepGC -XX:+UseParNewGC -XX:+CMSClassUnloadingEnabled -XX:-CMSParallelRemarkEnabled -XX:CMSInitiatingOccupancyFraction=80 -XX:SoftRefLRUPolicyMSPerMB=0 -Xloggc:$DORIS_HOME/log/fe.gc.log.$DATE -Dcom.mysql.cj.disableAbandonedConnectionCleanup=true" # For jdk 17+, this JAVA_OPTS will be used as default JVM options -JAVA_OPTS_FOR_JDK_17="-Dfile.encoding=UTF-8 -Djavax.security.auth.useSubjectCredsOnly=false -Xmx8192m -Xms8192m -XX:+UseG1GC -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=$LOG_DIR -Xlog:gc*,classhisto*=trace:$LOG_DIR/fe.gc.log.$CUR_DATE:time,uptime:filecount=10,filesize=50M -Darrow.enable_null_check_for_get=false --add-opens=java.base/java.lang=ALL-UNNAMED --add-opens=java.base/java.lang.invoke=ALL-UNNAMED --add-opens=java.base/java.lang.reflect=ALL-UNNAMED --add-opens=java.base/java.io=ALL-UNNAMED --add-opens=java.base/java.net=ALL-UNNAMED --add-opens=java.base/java.nio=ALL-UNNAMED --add-opens=java.base/java.util=ALL-UNNAMED --add-opens=java.base/java.util.concurrent=ALL-UNNAMED --add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED --add-opens=java.base/sun.nio.ch=ALL-UNNAMED --add-opens=java.base/sun.nio.cs=ALL-UNNAMED --add-opens=java.base/sun.security.action=ALL-UNNAMED --add-opens=java.base/sun.util.calendar=ALL-UNNAMED --add-opens=java.security.jgss/sun.security.krb5=ALL-UNNAMED --add-opens=java.management/sun.management=ALL-UNNAMED --add-opens=java.base/jdk.internal.ref=ALL-UNNAMED --add-opens=java.xml/com.sun.org.apache.xerces.internal.jaxp=ALL-UNNAMED" +JAVA_OPTS_FOR_JDK_17="-Dfile.encoding=UTF-8 -Djavax.security.auth.useSubjectCredsOnly=false -Xmx4096m -XX:+UseG1GC -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=$LOG_DIR -Xlog:gc*,classhisto*=trace:$LOG_DIR/fe.gc.log.$CUR_DATE:time,uptime:filecount=10,filesize=50M -Darrow.enable_null_check_for_get=false --add-opens=java.base/java.lang=ALL-UNNAMED --add-opens=java.base/java.lang.invoke=ALL-UNNAMED --add-opens=java.base/java.lang.reflect=ALL-UNNAMED --add-opens=java.base/java.io=ALL-UNNAMED --add-opens=java.base/java.net=ALL-UNNAMED --add-opens=java.base/java.nio=ALL-UNNAMED --add-opens=java.base/java.util=ALL-UNNAMED --add-opens=java.base/java.util.concurrent=ALL-UNNAMED --add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED --add-opens=java.base/sun.nio.ch=ALL-UNNAMED --add-opens=java.base/sun.nio.cs=ALL-UNNAMED --add-opens=java.base/sun.security.action=ALL-UNNAMED --add-opens=java.base/sun.util.calendar=ALL-UNNAMED --add-opens=java.security.jgss/sun.security.krb5=ALL-UNNAMED --add-opens=java.management/sun.management=ALL-UNNAMED --add-opens=java.base/jdk.internal.ref=ALL-UNNAMED --add-opens=java.xml/com.sun.org.apache.xerces.internal.jaxp=ALL-UNNAMED" ## ## the lowercase properties are read by main program.