From 5f4541c56983dc1da800f3b2bbd4887424eb740e Mon Sep 17 00:00:00 2001 From: cambyzhu Date: Fri, 14 Aug 2026 14:26:26 +0800 Subject: [PATCH 1/2] [feat](io) HDFS lazy open + iceberg delete file file_size propagation --- .../iceberg_delete_file_reader_helper.cpp | 9 +- .../table/iceberg_delete_file_reader_helper.h | 2 +- be/src/format_v2/table/iceberg_reader.cpp | 4 +- be/src/io/fs/file_handle_cache.cpp | 35 +++-- be/src/io/fs/file_handle_cache.h | 9 +- be/src/io/fs/hdfs_file_reader.cpp | 3 +- be/src/io/fs/s3_file_reader.cpp | 5 + ...iceberg_delete_file_reader_helper_test.cpp | 8 +- .../format_v2/table/iceberg_reader_test.cpp | 122 ++++++++++++++++++ be/test/io/fs/file_handle_cache_test.cpp | 28 ++++ .../iceberg/IcebergScanPlanProvider.java | 6 +- .../connector/iceberg/IcebergScanRange.java | 20 ++- .../iceberg/IcebergScanPlanProviderTest.java | 39 ++++++ .../iceberg/IcebergScanRangeTest.java | 40 ++++-- gensrc/thrift/PlanNodes.thrift | 1 + run-be-ut.sh | 7 + 16 files changed, 295 insertions(+), 43 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 f8828e402b03e9..12355493f0e908 100644 --- a/be/src/format/table/iceberg_delete_file_reader_helper.cpp +++ b/be/src/format/table/iceberg_delete_file_reader_helper.cpp @@ -236,12 +236,13 @@ TFileScanRangeParams build_iceberg_delete_scan_range_params( return params; } -TFileRangeDesc build_iceberg_delete_file_range(const std::string& path) { +TFileRangeDesc build_iceberg_delete_file_range(const std::string& path, int64_t file_size) { TFileRangeDesc range; range.path = path; range.start_offset = 0; range.size = -1; - range.file_size = -1; + // thrift optional defaults to 0 when unset; treat 0 as unknown (-1) + range.file_size = (file_size <= 0) ? -1 : file_size; return range; } @@ -276,7 +277,7 @@ Status read_iceberg_position_delete_file(const TIcebergDeleteFileDesc& delete_fi return Status::InvalidArgument("invalid position delete reader options"); } - TFileRangeDesc delete_range = build_iceberg_delete_file_range(delete_file.path); + TFileRangeDesc delete_range = build_iceberg_delete_file_range(delete_file.path, delete_file.file_size); if (options.fs_name != nullptr && !options.fs_name->empty()) { delete_range.__set_fs_name(*options.fs_name); } @@ -348,7 +349,7 @@ Status read_iceberg_deletion_vector(const TIcebergDeleteFileDesc& delete_file, 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); + TFileRangeDesc delete_range = build_iceberg_delete_file_range(delete_file.path, delete_file.file_size); if (options.fs_name != nullptr && !options.fs_name->empty()) { delete_range.__set_fs_name(*options.fs_name); } diff --git a/be/src/format/table/iceberg_delete_file_reader_helper.h b/be/src/format/table/iceberg_delete_file_reader_helper.h index adc0ef196f4b75..d43443ecd8e1c9 100644 --- a/be/src/format/table/iceberg_delete_file_reader_helper.h +++ b/be/src/format/table/iceberg_delete_file_reader_helper.h @@ -67,7 +67,7 @@ TFileScanRangeParams build_iceberg_delete_scan_range_params( const std::map& hadoop_conf, TFileType::type file_type, const std::vector& broker_addresses); -TFileRangeDesc build_iceberg_delete_file_range(const std::string& path); +TFileRangeDesc build_iceberg_delete_file_range(const std::string& path, int64_t file_size); bool is_iceberg_deletion_vector(const TIcebergDeleteFileDesc& delete_file); diff --git a/be/src/format_v2/table/iceberg_reader.cpp b/be/src/format_v2/table/iceberg_reader.cpp index 37ed4c00e28a56..895ba8d9b09c4f 100644 --- a/be/src/format_v2/table/iceberg_reader.cpp +++ b/be/src/format_v2/table/iceberg_reader.cpp @@ -1074,7 +1074,7 @@ Status IcebergTableReader::_parse_deletion_vector_file(const TTableFormatFileDes desc->path = deletion_vector->path; desc->start_offset = deletion_vector->content_offset; desc->size = static_cast(bytes_read); - desc->file_size = -1; + desc->file_size = deletion_vector->file_size; desc->format = DeleteFileDesc::Format::ICEBERG; *has_delete_file = true; return Status::OK(); @@ -1430,7 +1430,7 @@ Status IcebergTableReader::_create_delete_file_reader(const TIcebergDeleteFileDe return Status::NotSupported("Unsupported Iceberg delete file format {}", delete_file.file_format); } - auto delete_range = build_iceberg_delete_file_range(delete_file.path); + auto delete_range = build_iceberg_delete_file_range(delete_file.path, delete_file.file_size); if (_current_task != nullptr && _current_task->data_file != nullptr && !_current_task->data_file->fs_name.empty()) { delete_range.__set_fs_name(_current_task->data_file->fs_name); diff --git a/be/src/io/fs/file_handle_cache.cpp b/be/src/io/fs/file_handle_cache.cpp index 41617ba10159fc..6153c2b62c87d2 100644 --- a/be/src/io/fs/file_handle_cache.cpp +++ b/be/src/io/fs/file_handle_cache.cpp @@ -41,19 +41,9 @@ HdfsFileHandle::~HdfsFileHandle() { } Status HdfsFileHandle::init(int64_t file_size) { - _hdfs_file = hdfsOpenFile(_fs, _fname.c_str(), O_RDONLY, 0, 0, 0); - if (_hdfs_file == nullptr) { - std::string _err_msg = hdfs_error(); - // invoker maybe just skip Status.NotFound and continue - // so we need distinguish between it and other kinds of errors - if (_err_msg.find("No such file or directory") != std::string::npos) { - return Status::NotFound(_err_msg); - } - return Status::InternalError("failed to open {}: {}", _fname, _err_msg); - } - _file_size = file_size; if (_file_size <= 0) { + // file_size unknown, fetch via hdfsGetPathInfo (no need to hdfsOpenFile) hdfsFileInfo* file_info = hdfsGetPathInfo(_fs, _fname.c_str()); if (file_info == nullptr) { return Status::InternalError("failed to get file size of {}: {}", _fname, hdfs_error()); @@ -64,6 +54,21 @@ Status HdfsFileHandle::init(int64_t file_size) { return Status::OK(); } +Status HdfsFileHandle::ensure_open() { + std::call_once(_open_once, [this]() { + VLOG_DEBUG << "lazy open hdfs file: " << _fname; + _hdfs_file = hdfsOpenFile(_fs, _fname.c_str(), O_RDONLY, 0, 0, 0); + }); + if (_hdfs_file == nullptr) { + std::string _err_msg = hdfs_error(); + if (_err_msg.find("No such file or directory") != std::string::npos) { + return Status::NotFound(_err_msg); + } + return Status::InternalError("failed to open {}: {}", _fname, _err_msg); + } + return Status::OK(); +} + CachedHdfsFileHandle::CachedHdfsFileHandle(const hdfsFS& fs, const std::string& fname, int64_t mtime) : HdfsFileHandle(fs, fname, mtime) {} @@ -98,8 +103,14 @@ void FileHandleCache::Accessor::destroy() { FileHandleCache::Accessor::~Accessor() { if (_cache_accessor.get()) { + auto* handle = get(); + if (handle->file() == nullptr) { + // Not opened (lazy open), no resources to release + release(); + return; + } #ifdef USE_HADOOP_HDFS - if (hdfsUnbufferFile(get()->file()) != 0) { + if (hdfsUnbufferFile(handle->file()) != 0) { VLOG_FILE << "FS does not support file handle unbuffering, closing file=" << _cache_accessor.get_key()->second.first; destroy(); diff --git a/be/src/io/fs/file_handle_cache.h b/be/src/io/fs/file_handle_cache.h index ce3c708ba99242..2648ccd345d58f 100644 --- a/be/src/io/fs/file_handle_cache.h +++ b/be/src/io/fs/file_handle_cache.h @@ -26,6 +26,7 @@ #include #include #include +#include #include #include @@ -49,9 +50,12 @@ class HdfsFileHandle { /// Destructor will close the file handle ~HdfsFileHandle(); - /// Init opens the file handle + /// Init only sets file_size (from param or hdfsGetPathInfo), does NOT open the file. Status init(int64_t file_size); + /// Lazily opens the file handle on first read. Thread-safe via std::call_once. + Status ensure_open(); + hdfsFS fs() const { return _fs; } hdfsFile file() const { return _hdfs_file; } int64_t mtime() const { return _mtime; } @@ -66,7 +70,8 @@ class HdfsFileHandle { const std::string _fname; hdfsFile _hdfs_file = nullptr; int64_t _mtime; - int64_t _file_size; + int64_t _file_size = -1; + std::once_flag _open_once; }; /// CachedHdfsFileHandles are owned by the file handle cache and are used for no diff --git a/be/src/io/fs/hdfs_file_reader.cpp b/be/src/io/fs/hdfs_file_reader.cpp index 2f734c9356b663..ee439ea319ccef 100644 --- a/be/src/io/fs/hdfs_file_reader.cpp +++ b/be/src/io/fs/hdfs_file_reader.cpp @@ -123,6 +123,7 @@ Status HdfsFileReader::close() { Status HdfsFileReader::read_at_impl(size_t offset, Slice result, size_t* bytes_read, const IOContext* io_ctx) { SCOPED_TIMER(_total_read_time); + RETURN_IF_ERROR(_handle->ensure_open()); auto st = do_read_at_impl(offset, result, bytes_read, io_ctx); if (!st.ok()) { _handle = nullptr; @@ -257,7 +258,7 @@ Status HdfsFileReader::do_read_at_impl(size_t offset, Slice result, size_t* byte void HdfsFileReader::_collect_profile_before_close() { if (_profile != nullptr && is_hdfs(_fs_name)) { #ifdef USE_HADOOP_HDFS - if (_handle == nullptr) [[unlikely]] { + if (_handle == nullptr || _handle->file() == nullptr) [[unlikely]] { return; } diff --git a/be/src/io/fs/s3_file_reader.cpp b/be/src/io/fs/s3_file_reader.cpp index 8a6e5c0fdc4978..6d6c4af9723384 100644 --- a/be/src/io/fs/s3_file_reader.cpp +++ b/be/src/io/fs/s3_file_reader.cpp @@ -68,12 +68,17 @@ Result S3FileReader::create(std::shared_ptr()) << oversized_status; + EXPECT_TRUE(oversized_status.template is()) << oversized_status; EXPECT_NE(oversized_status.to_string().find("range exceeds file size"), std::string::npos); EXPECT_NE(oversized_status.to_string().find(dv_path), std::string::npos); } @@ -613,6 +613,8 @@ TEST(IcebergDeleteFileReaderHelperTest, ReadMixedEncodingParquetPositionDeleteFi delete_file.path = kMixedPositionDeleteFile; delete_file.file_format = TFileFormatType::FORMAT_PARQUET; delete_file.__isset.file_format = true; + delete_file.__isset.file_size = true; + delete_file.file_size = -1; IcebergDeleteFileReaderOptions options; options.state = &runtime_state; diff --git a/be/test/format_v2/table/iceberg_reader_test.cpp b/be/test/format_v2/table/iceberg_reader_test.cpp index 5c78f6b19efad3..078081dc21c0f7 100644 --- a/be/test/format_v2/table/iceberg_reader_test.cpp +++ b/be/test/format_v2/table/iceberg_reader_test.cpp @@ -1208,6 +1208,8 @@ TIcebergDeleteFileDesc make_iceberg_equality_delete_file( delete_file.__set_path(path); delete_file.__set_field_ids(field_ids); delete_file.__set_file_format(file_format); + // Set file_size to actual file size, simulating FE propagation from iceberg manifest + delete_file.__set_file_size(static_cast(std::filesystem::file_size(path))); return delete_file; } @@ -4923,5 +4925,125 @@ TEST(IcebergV2ReaderTest, DataFileIsMarkedImmutableForPageCache) { EXPECT_TRUE(reader.current_data_file_is_immutable()); } +// E2E: Verify file_size from FE thrift propagates through the full reader chain. +// Data file: 3 rows (id=1,2,3). Equality delete: delete id=2. +// With file_size set on TIcebergDeleteFileDesc, the reader should use it +// instead of falling back to stat, and return correct results (id=1,3). +TEST(IcebergV2ReaderTest, IcebergEqualityDeleteFileSizePropagatedToReader) { + const auto test_dir = + std::filesystem::temp_directory_path() / "doris_iceberg_eq_delete_file_size_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + + const auto file_path = (test_dir / "split.parquet").string(); + const auto delete_file_path = (test_dir / "equality-delete.parquet").string(); + write_int_pair_parquet_file(file_path, {1, 2, 3}, {10, 20, 30}, {"one", "two", "three"}); + write_iceberg_equality_delete_parquet_file(delete_file_path, 0, 2); + + const auto delete_file_size = static_cast(std::filesystem::file_size(delete_file_path)); + ASSERT_GT(delete_file_size, 0) << "delete file should not be empty"; + + std::vector projected_columns; + projected_columns.push_back(make_table_column(0, "id", std::make_shared())); + + 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 = nullptr, + }) + .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_equality_delete_file(delete_file_path, {0})})); + ASSERT_TRUE(reader.prepare_split(split_options).ok()); + + Block block = build_table_block(projected_columns); + bool eos = false; + ASSERT_TRUE(reader.get_block(&block, &eos).ok()); + ASSERT_FALSE(eos); + ASSERT_EQ(block.rows(), 2); + const auto& id_column = assert_cast(expect_not_null_table_column(block, 0)); + EXPECT_EQ(id_column.get_element(0), 1); + EXPECT_EQ(id_column.get_element(1), 3); + + ASSERT_TRUE(reader.close().ok()); + std::filesystem::remove_all(test_dir); +} + +// E2E: Verify that when file_size is not set (thrift optional defaults to 0), +// build_iceberg_delete_file_range converts 0 to -1, and FileFactory falls back +// to stat. The reader should still return correct results. +TEST(IcebergV2ReaderTest, IcebergEqualityDeleteFileSizeUnknownFallsBackToStat) { + const auto test_dir = + std::filesystem::temp_directory_path() / "doris_iceberg_eq_delete_no_file_size_test"; + std::filesystem::remove_all(test_dir); + std::filesystem::create_directories(test_dir); + + const auto file_path = (test_dir / "split.parquet").string(); + const auto delete_file_path = (test_dir / "equality-delete.parquet").string(); + write_int_pair_parquet_file(file_path, {1, 2, 3}, {10, 20, 30}, {"one", "two", "three"}); + write_iceberg_equality_delete_parquet_file(delete_file_path, 0, 2); + + // Construct delete file WITHOUT file_size (simulating FE not setting it) + TIcebergDeleteFileDesc delete_file; + delete_file.__set_content(2); + delete_file.__set_path(delete_file_path); + delete_file.__set_field_ids({0}); + delete_file.__set_file_format(TFileFormatType::FORMAT_PARQUET); + // file_size intentionally not set — thrift defaults to 0 + + std::vector projected_columns; + projected_columns.push_back(make_table_column(0, "id", std::make_shared())); + + 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 = nullptr, + }) + .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, {delete_file})); + ASSERT_TRUE(reader.prepare_split(split_options).ok()); + + Block block = build_table_block(projected_columns); + bool eos = false; + ASSERT_TRUE(reader.get_block(&block, &eos).ok()); + ASSERT_FALSE(eos); + ASSERT_EQ(block.rows(), 2); + const auto& id_column = assert_cast(expect_not_null_table_column(block, 0)); + EXPECT_EQ(id_column.get_element(0), 1); + EXPECT_EQ(id_column.get_element(1), 3); + + ASSERT_TRUE(reader.close().ok()); + std::filesystem::remove_all(test_dir); +} + } // namespace } // namespace doris::format diff --git a/be/test/io/fs/file_handle_cache_test.cpp b/be/test/io/fs/file_handle_cache_test.cpp index 5c1f7d1d9e05e8..41b503a49cd4bf 100644 --- a/be/test/io/fs/file_handle_cache_test.cpp +++ b/be/test/io/fs/file_handle_cache_test.cpp @@ -22,6 +22,8 @@ #include #include +#include "format/table/iceberg_delete_file_reader_helper.h" + namespace doris::io { TEST(FileHandleCacheTest, CacheKeyIncludesHdfsFs) { @@ -40,4 +42,30 @@ TEST(FileHandleCacheTest, CacheKeyIncludesHdfsFs) { mtime + 1)); } +// Verify that init() with file_size > 0 does NOT open the file (lazy open). +// The handle should know file_size but file() should be nullptr. +TEST(FileHandleCacheTest, InitWithKnownFileSizeDoesNotOpenFile) { + auto mock_fs = reinterpret_cast(static_cast(0x1)); + ExclusiveHdfsFileHandle handle(mock_fs, "/nonexistent/file.parquet", 12345); + auto st = handle.init(4096); + ASSERT_TRUE(st.ok()) << st; + // file_size should be set from parameter + EXPECT_EQ(handle.file_size(), 4096); + // file() should be nullptr — lazy open means no hdfsOpenFile was called + EXPECT_EQ(handle.file(), nullptr); +} + +// Verify that build_iceberg_delete_file_range treats file_size <= 0 as unknown (-1). +// This covers the case where thrift optional file_size defaults to 0. +TEST(FileHandleCacheTest, BuildDeleteFileRangeTreatsZeroAsUnknown) { + auto range_known = build_iceberg_delete_file_range("s3://b/f.parquet", 1024); + EXPECT_EQ(range_known.file_size, 1024); + + auto range_zero = build_iceberg_delete_file_range("s3://b/f.parquet", 0); + EXPECT_EQ(range_zero.file_size, -1); + + auto range_neg = build_iceberg_delete_file_range("s3://b/f.parquet", -1); + EXPECT_EQ(range_neg.file_size, -1); +} + } // namespace doris::io diff --git a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanPlanProvider.java b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanPlanProvider.java index 02c5ac5bcf8f53..d7bfbd9a1bb432 100644 --- a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanPlanProvider.java +++ b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanPlanProvider.java @@ -1481,13 +1481,13 @@ IcebergScanRange.DeleteFile convertDelete(DeleteFile delete, UnaryOperator fieldIds, Long contentOffset, Long contentSizeInBytes) { + Long positionUpperBound, List fieldIds, Long contentOffset, Long contentSizeInBytes, + long fileSize) { this.path = path; this.content = content; this.fileFormat = fileFormat; @@ -642,13 +644,14 @@ private DeleteFile(String path, int content, TFileFormatType fileFormat, Long po this.fieldIds = fieldIds != null ? Collections.unmodifiableList(new ArrayList<>(fieldIds)) : null; this.contentOffset = contentOffset; this.contentSizeInBytes = contentSizeInBytes; + this.fileSize = fileSize; } /** A position delete file (content 1): row positions to drop, with optional [lower,upper] bounds. */ public static DeleteFile positionDelete(String path, TFileFormatType fileFormat, - Long positionLowerBound, Long positionUpperBound) { + Long positionLowerBound, Long positionUpperBound, long fileSize) { return new DeleteFile(path, CONTENT_POSITION_DELETE, fileFormat, - positionLowerBound, positionUpperBound, null, null, null); + positionLowerBound, positionUpperBound, null, null, null, fileSize); } /** @@ -657,14 +660,16 @@ public static DeleteFile positionDelete(String path, TFileFormatType fileFormat, * bounds (legacy {@code DeletionVector extends PositionDelete}); {@code file_format} stays unset. */ public static DeleteFile deletionVector(String path, Long positionLowerBound, Long positionUpperBound, - long contentOffset, long contentSizeInBytes) { + long contentOffset, long contentSizeInBytes, long fileSize) { return new DeleteFile(path, CONTENT_DELETION_VECTOR, null, - positionLowerBound, positionUpperBound, null, contentOffset, contentSizeInBytes); + positionLowerBound, positionUpperBound, null, contentOffset, contentSizeInBytes, fileSize); } /** An equality delete file (content 2): rows equal on {@code fieldIds} are dropped (BE re-projects). */ - public static DeleteFile equalityDelete(String path, TFileFormatType fileFormat, List fieldIds) { - return new DeleteFile(path, CONTENT_EQUALITY_DELETE, fileFormat, null, null, fieldIds, null, null); + public static DeleteFile equalityDelete(String path, TFileFormatType fileFormat, List fieldIds, + long fileSize) { + return new DeleteFile(path, CONTENT_EQUALITY_DELETE, fileFormat, null, null, + fieldIds, null, null, fileSize); } int getContent() { @@ -693,6 +698,7 @@ TIcebergDeleteFileDesc toThrift() { desc.setContentSizeInBytes(contentSizeInBytes); } desc.setContent(content); + desc.setFileSize(fileSize); return desc; } } diff --git a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanPlanProviderTest.java b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanPlanProviderTest.java index 911b6073a9b091..0d7c444af8763c 100644 --- a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanPlanProviderTest.java +++ b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanPlanProviderTest.java @@ -3913,4 +3913,43 @@ public void deleteFile(String path) { throw new UnsupportedOperationException(); } } + + @Test + public void deleteFileSizePropagatedFromIcebergManifest() { + // E2E: iceberg DeleteFile.fileSizeInBytes (482) → IcebergScanPlanProvider → + // IcebergScanRange.DeleteFile.fileSize → toThrift → TIcebergDeleteFileDesc.file_size + Table table = createTable("t_filesize", SCHEMA, PartitionSpec.unpartitioned(), + Collections.singletonMap("format-version", "2")); + table.newAppend() + .appendFile(dataFile(table.spec(), "s3://b/db/t_filesize/f1.parquet", 512, null, null)) + .commit(); + // Position delete with fileSizeInBytes=128 (rewritableDeleteDescs includes position deletes) + DeleteFile posDelete = FileMetadata.deleteFileBuilder(table.spec()) + .ofPositionDeletes() + .withPath("s3://b/db/t_filesize/pos-delete.parquet") + .withFormat(FileFormat.PARQUET) + .withFileSizeInBytes(128L) + .withRecordCount(1L) + .build(); + table.newRowDelta().addDeletes(posDelete).commit(); + + IcebergScanPlanProvider provider = new IcebergScanPlanProvider( + IcebergCatalogProperties.of(Collections.emptyMap()), opsReturning(table)); + List ranges = provider.planScan( + new FakeScanSession("UTC", Collections.emptyMap()), + ConnectorScanRequest.builder(new IcebergTableHandle("db1", "t_filesize"), Collections.emptyList()) + .build()); + + Assertions.assertEquals(1, ranges.size()); + IcebergScanRange range = (IcebergScanRange) ranges.get(0); + + // Use rewritableDeleteDescs() (returns TIcebergDeleteFileDesc via toThrift) + // to verify file_size propagated from iceberg manifest → DeleteFile → toThrift. + // Note: equality deletes are excluded from rewritableDeleteDescs, so use position delete instead. + List descs = range.rewritableDeleteDescs(); + Assertions.assertFalse(descs.isEmpty()); + TIcebergDeleteFileDesc desc = descs.get(0); + Assertions.assertTrue(desc.isSetFileSize()); + Assertions.assertEquals(128L, desc.getFileSize()); + } } diff --git a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanRangeTest.java b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanRangeTest.java index 07f52b5d3a50a6..01604230db9c10 100644 --- a/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanRangeTest.java +++ b/fe/fe-connector/fe-connector-iceberg/src/test/java/org/apache/doris/connector/iceberg/IcebergScanRangeTest.java @@ -191,7 +191,7 @@ public void populateRangeParamsV2EmitsPositionDeleteFile() { // A position delete (content 1) with parquet format + [lower,upper] bounds. MUTATION: dropping the // bounds, wrong content id, or wrong format -> red. IcebergScanRange.DeleteFile posDelete = IcebergScanRange.DeleteFile.positionDelete( - "s3://b/db/t/pos-delete.parquet", TFileFormatType.FORMAT_PARQUET, 10L, 99L); + "s3://b/db/t/pos-delete.parquet", TFileFormatType.FORMAT_PARQUET, 10L, 99L, 100L); IcebergScanRange range = new IcebergScanRange.Builder() .path("s3://b/db/t/f.parquet").fileFormat("parquet").formatVersion(2) .deleteFiles(Collections.singletonList(posDelete)).build(); @@ -218,7 +218,7 @@ public void populateRangeParamsV2EmitsPositionDeleteWithoutBounds() { // No bounds present -> position_lower/upper_bound left UNSET (legacy emits them only when present). // MUTATION: defaulting an absent bound to 0 / -1 instead of unset -> red. IcebergScanRange.DeleteFile posDelete = IcebergScanRange.DeleteFile.positionDelete( - "s3://b/db/t/pos-delete.orc", TFileFormatType.FORMAT_ORC, null, null); + "s3://b/db/t/pos-delete.orc", TFileFormatType.FORMAT_ORC, null, null, 100L); IcebergScanRange range = new IcebergScanRange.Builder() .path("s3://b/db/t/f.parquet").fileFormat("parquet").formatVersion(2) .deleteFiles(Collections.singletonList(posDelete)).build(); @@ -236,9 +236,9 @@ public void populateRangeParamsV2EmitsDeletionVectorAndEqualityDelete() { // A deletion vector (content 3, PUFFIN): blob content_offset/size set, file_format UNSET, bounds // carried (it IS a position delete). An equality delete (content 2): field-ids set, no bounds/blob. IcebergScanRange.DeleteFile dv = IcebergScanRange.DeleteFile.deletionVector( - "s3://b/db/t/dv.puffin", 5L, 42L, 16L, 64L); + "s3://b/db/t/dv.puffin", 5L, 42L, 16L, 64L, 100L); IcebergScanRange.DeleteFile eq = IcebergScanRange.DeleteFile.equalityDelete( - "s3://b/db/t/eq-delete.parquet", TFileFormatType.FORMAT_PARQUET, Arrays.asList(3, 7)); + "s3://b/db/t/eq-delete.parquet", TFileFormatType.FORMAT_PARQUET, Arrays.asList(3, 7), 100L); IcebergScanRange range = new IcebergScanRange.Builder() .path("s3://b/db/t/f.parquet").fileFormat("parquet").formatVersion(2) .deleteFiles(Arrays.asList(dv, eq)).build(); @@ -467,11 +467,11 @@ public void rewritableDeleteDescsKeepsDvAndPositionButDropsEquality() { // rewritten, so they MUST be excluded (mirrors legacy deleteFilesDescByReferencedDataFile). MUTATION: // including the equality delete -> the BE would treat equality rows as positions / over-delete. IcebergScanRange.DeleteFile dv = IcebergScanRange.DeleteFile.deletionVector( - "s3://b/db/t/dv.puffin", 5L, 42L, 16L, 64L); + "s3://b/db/t/dv.puffin", 5L, 42L, 16L, 64L, 100L); IcebergScanRange.DeleteFile pos = IcebergScanRange.DeleteFile.positionDelete( - "s3://b/db/t/pos.parquet", TFileFormatType.FORMAT_PARQUET, 1L, 9L); + "s3://b/db/t/pos.parquet", TFileFormatType.FORMAT_PARQUET, 1L, 9L, 100L); IcebergScanRange.DeleteFile eq = IcebergScanRange.DeleteFile.equalityDelete( - "s3://b/db/t/eq.parquet", TFileFormatType.FORMAT_PARQUET, Arrays.asList(3, 7)); + "s3://b/db/t/eq.parquet", TFileFormatType.FORMAT_PARQUET, Arrays.asList(3, 7), 100L); IcebergScanRange range = new IcebergScanRange.Builder() .path("s3://b/db/t/f.parquet").fileFormat("parquet").formatVersion(3) .deleteFiles(Arrays.asList(dv, pos, eq)).build(); @@ -495,10 +495,34 @@ public void rewritableDeleteDescsEmptyWhenNoDeletesOrAllEquality() { Assertions.assertTrue(none.rewritableDeleteDescs().isEmpty()); IcebergScanRange.DeleteFile eq = IcebergScanRange.DeleteFile.equalityDelete( - "s3://b/db/t/eq.parquet", TFileFormatType.FORMAT_PARQUET, Arrays.asList(3)); + "s3://b/db/t/eq.parquet", TFileFormatType.FORMAT_PARQUET, Arrays.asList(3), 100L); IcebergScanRange onlyEq = new IcebergScanRange.Builder() .path("s3://b/db/t/f.parquet").fileFormat("parquet").formatVersion(3) .deleteFiles(Collections.singletonList(eq)).build(); Assertions.assertTrue(onlyEq.rewritableDeleteDescs().isEmpty()); } + + @Test + public void deleteFileToThriftPropagatesFileSize() { + // Equality delete + IcebergScanRange.DeleteFile eq = IcebergScanRange.DeleteFile.equalityDelete( + "s3://b/db/t/eq.parquet", TFileFormatType.FORMAT_PARQUET, Arrays.asList(1), 482L); + TIcebergDeleteFileDesc eqDesc = eq.toThrift(); + Assertions.assertTrue(eqDesc.isSetFileSize()); + Assertions.assertEquals(482L, eqDesc.getFileSize()); + + // Position delete + IcebergScanRange.DeleteFile pos = IcebergScanRange.DeleteFile.positionDelete( + "s3://b/db/t/pos.parquet", TFileFormatType.FORMAT_PARQUET, 10L, 99L, 735L); + TIcebergDeleteFileDesc posDesc = pos.toThrift(); + Assertions.assertTrue(posDesc.isSetFileSize()); + Assertions.assertEquals(735L, posDesc.getFileSize()); + + // Deletion vector + IcebergScanRange.DeleteFile dv = IcebergScanRange.DeleteFile.deletionVector( + "s3://b/db/t/dv.puffin", 5L, 42L, 16L, 64L, 1000L); + TIcebergDeleteFileDesc dvDesc = dv.toThrift(); + Assertions.assertTrue(dvDesc.isSetFileSize()); + Assertions.assertEquals(1000L, dvDesc.getFileSize()); + } } diff --git a/gensrc/thrift/PlanNodes.thrift b/gensrc/thrift/PlanNodes.thrift index c71c9e0a3b0e55..cf2bc594405a0c 100644 --- a/gensrc/thrift/PlanNodes.thrift +++ b/gensrc/thrift/PlanNodes.thrift @@ -328,6 +328,7 @@ struct TIcebergDeleteFileDesc { 9: optional string original_path; // Referenced data file path. Required to materialize rows from deletion vectors. 10: optional string referenced_data_file_path; + 11: optional i64 file_size; } struct TIcebergFileDesc { diff --git a/run-be-ut.sh b/run-be-ut.sh index 3ff9cd573ee1ce..7a1d4c0ca71664 100755 --- a/run-be-ut.sh +++ b/run-be-ut.sh @@ -423,6 +423,13 @@ if [[ -d "${DORIS_TEST_BINARY_DIR}/util/test_data" ]]; then fi cp -r "${DORIS_HOME}/be/test/util/test_data" "${DORIS_TEST_BINARY_DIR}/util"/ +# prepare exec test_data (iceberg delete file tests use ./be/test/exec/test_data/...) +mkdir -p "${DORIS_TEST_BINARY_DIR}/be/test/exec" +if [[ -d "${DORIS_TEST_BINARY_DIR}/be/test/exec/test_data" ]]; then + rm -rf "${DORIS_TEST_BINARY_DIR}/be/test/exec/test_data" +fi +cp -r "${DORIS_HOME}/be/test/exec/test_data" "${DORIS_TEST_BINARY_DIR}/be/test/exec"/ + # prepare ut temp dir UT_TMP_DIR="${DORIS_HOME}/ut_dir" rm -rf "${UT_TMP_DIR}" From 4826368125ccf52ea517a413d55d4347301d1a57 Mon Sep 17 00:00:00 2001 From: cambyzhu Date: Fri, 14 Aug 2026 15:52:13 +0800 Subject: [PATCH 2/2] [chore](io) fix clang-format for lazy open changes --- be/src/format/table/iceberg_delete_file_reader_helper.cpp | 6 ++++-- be/src/io/fs/s3_file_reader.cpp | 5 +++-- .../iceberg/iceberg_delete_file_reader_helper_test.cpp | 3 ++- be/test/format_v2/table/iceberg_reader_test.cpp | 3 ++- 4 files changed, 11 insertions(+), 6 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 12355493f0e908..52369e3271dcb6 100644 --- a/be/src/format/table/iceberg_delete_file_reader_helper.cpp +++ b/be/src/format/table/iceberg_delete_file_reader_helper.cpp @@ -277,7 +277,8 @@ Status read_iceberg_position_delete_file(const TIcebergDeleteFileDesc& delete_fi return Status::InvalidArgument("invalid position delete reader options"); } - TFileRangeDesc delete_range = build_iceberg_delete_file_range(delete_file.path, delete_file.file_size); + TFileRangeDesc delete_range = + build_iceberg_delete_file_range(delete_file.path, delete_file.file_size); if (options.fs_name != nullptr && !options.fs_name->empty()) { delete_range.__set_fs_name(*options.fs_name); } @@ -349,7 +350,8 @@ Status read_iceberg_deletion_vector(const TIcebergDeleteFileDesc& delete_file, 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, delete_file.file_size); + TFileRangeDesc delete_range = + build_iceberg_delete_file_range(delete_file.path, delete_file.file_size); if (options.fs_name != nullptr && !options.fs_name->empty()) { delete_range.__set_fs_name(*options.fs_name); } diff --git a/be/src/io/fs/s3_file_reader.cpp b/be/src/io/fs/s3_file_reader.cpp index 6d6c4af9723384..5de990e82c34be 100644 --- a/be/src/io/fs/s3_file_reader.cpp +++ b/be/src/io/fs/s3_file_reader.cpp @@ -68,8 +68,9 @@ Result S3FileReader::create(std::shared_ptr(std::filesystem::file_size(delete_file_path)); + const auto delete_file_size = + static_cast(std::filesystem::file_size(delete_file_path)); ASSERT_GT(delete_file_size, 0) << "delete file should not be empty"; std::vector projected_columns;