Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 7 additions & 4 deletions be/src/format/table/iceberg_delete_file_reader_helper.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

Expand Down Expand Up @@ -276,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);
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);
}
Expand Down Expand Up @@ -348,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);
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);
}
Expand Down
2 changes: 1 addition & 1 deletion be/src/format/table/iceberg_delete_file_reader_helper.h
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ TFileScanRangeParams build_iceberg_delete_scan_range_params(
const std::map<std::string, std::string>& hadoop_conf, TFileType::type file_type,
const std::vector<TNetworkAddress>& 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);

Expand Down
4 changes: 2 additions & 2 deletions be/src/format_v2/table/iceberg_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<int64_t>(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();
Expand Down Expand Up @@ -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);
Expand Down
35 changes: 23 additions & 12 deletions be/src/io/fs/file_handle_cache.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand All @@ -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) {}
Expand Down Expand Up @@ -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();
Expand Down
9 changes: 7 additions & 2 deletions be/src/io/fs/file_handle_cache.h
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
#include <list>
#include <map>
#include <memory>
#include <mutex>
#include <string>
#include <utility>

Expand All @@ -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; }
Expand All @@ -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
Expand Down
3 changes: 2 additions & 1 deletion be/src/io/fs/hdfs_file_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}

Expand Down
6 changes: 6 additions & 0 deletions be/src/io/fs/s3_file_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -68,12 +68,18 @@ Result<FileReaderSPtr> S3FileReader::create(std::shared_ptr<const ObjClientHolde
std::string bucket, std::string key, int64_t file_size,
RuntimeProfile* profile) {
if (file_size < 0) {
VLOG_DEBUG
<< "S3FileReader: file_size not provided by FE, falling back to HeadObject: bucket="
<< bucket << ", key=" << key;
auto res = client->object_file_size(bucket, key);
if (!res.has_value()) {
return ResultError(std::move(res.error()));
}

file_size = res.value();
} else {
VLOG_DEBUG << "S3FileReader: file_size provided by FE: bucket=" << bucket << ", key=" << key
<< ", file_size=" << file_size;
}

return std::make_shared<S3FileReader>(std::move(client), std::move(bucket), std::move(key),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -214,7 +214,7 @@ IcebergDeleteFileReaderOptions delete_reader_options(RuntimeState* runtime_state
} // namespace

TEST(IcebergDeleteFileReaderHelperTest, BuildDeleteFileRange) {
auto range = build_iceberg_delete_file_range("s3://bucket/delete.parquet");
auto range = build_iceberg_delete_file_range("s3://bucket/delete.parquet", -1);
EXPECT_EQ(range.path, "s3://bucket/delete.parquet");
EXPECT_EQ(range.start_offset, 0);
EXPECT_EQ(range.size, -1);
Expand Down Expand Up @@ -384,7 +384,7 @@ TEST(IcebergDeleteFileReaderHelperTest, DeletionVectorReaderValidatesOpenedFileR
IcebergDeleteFileIOContext io_context(&state);

{
TFileRangeDesc exact_range = build_iceberg_delete_file_range(dv_path);
TFileRangeDesc exact_range = build_iceberg_delete_file_range(dv_path, -1);
exact_range.start_offset = 4;
exact_range.size = dv_size - exact_range.start_offset;
DeletionVectorReader exact_reader(&state, &profile, scan_params, exact_range,
Expand All @@ -397,7 +397,8 @@ TEST(IcebergDeleteFileReaderHelperTest, DeletionVectorReaderValidatesOpenedFileR
DeletionVectorReader oversized_reader(&state, &profile, scan_params, oversized_range,
&io_context.io_ctx);
const auto oversized_status = oversized_reader.open();
EXPECT_TRUE(oversized_status.is<ErrorCode::DATA_QUALITY_ERROR>()) << oversized_status;
EXPECT_TRUE(oversized_status.template is<ErrorCode::DATA_QUALITY_ERROR>())
<< 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);
}
Expand Down Expand Up @@ -613,6 +614,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;
Expand Down
123 changes: 123 additions & 0 deletions be/test/format_v2/table/iceberg_reader_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<int64_t>(std::filesystem::file_size(path)));
return delete_file;
}

Expand Down Expand Up @@ -4923,5 +4925,126 @@ 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<int64_t>(std::filesystem::file_size(delete_file_path));
ASSERT_GT(delete_file_size, 0) << "delete file should not be empty";

std::vector<ColumnDefinition> projected_columns;
projected_columns.push_back(make_table_column(0, "id", std::make_shared<DataTypeInt32>()));

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<const ColumnInt32&>(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<ColumnDefinition> projected_columns;
projected_columns.push_back(make_table_column(0, "id", std::make_shared<DataTypeInt32>()));

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<const ColumnInt32&>(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
Loading
Loading