This is an automated email from the ASF dual-hosted git repository.
morningman 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 9191b8e37e7 [feat](iceberg) HDFS lazy open + iceberg delete file
file_size propagation (#66773)
9191b8e37e7 is described below
commit 9191b8e37e7cf40c35b332cec40bbfa41b474fd4
Author: camby <[email protected]>
AuthorDate: Thu Sep 10 21:45:31 2026 +0800
[feat](iceberg) HDFS lazy open + iceberg delete file file_size propagation
(#66773)
### What problem does this PR solve?
Problem Summary:
When block cache hits, HdfsFileReader still calls hdfsOpenFile on every
construction, wasting NameNode RPCs for data that's never read from
HDFS.
How To Fix:
1. HdfsFileHandle now lazily opens on first read;
2. FE propagates delete file file_size through thrift to BE;
---
be/src/exec/scan/file_scanner.cpp | 10 +-
.../table/iceberg_delete_file_reader_helper.cpp | 10 +-
.../table/iceberg_delete_file_reader_helper.h | 2 +-
be/src/format/table/iceberg_reader_mixin.h | 9 +-
be/src/format_v2/orc/orc_reader.cpp | 9 +-
be/src/format_v2/table/iceberg_reader.cpp | 5 +-
be/src/io/fs/file_handle_cache.cpp | 68 +++-
be/src/io/fs/file_handle_cache.h | 10 +-
be/src/io/fs/hdfs_file_reader.cpp | 47 ++-
be/src/io/fs/hdfs_file_system.cpp | 9 +-
be/src/io/fs/local_file_reader.cpp | 2 +-
be/src/io/hdfs_util.cpp | 1 +
be/src/io/hdfs_util.h | 1 +
be/src/storage/index/index_file_reader.cpp | 10 +
.../index/inverted/inverted_index_fs_directory.cpp | 6 +-
be/src/storage/index/snii/snii_blob_directory.cpp | 6 +-
be/test/exec/scan/vfile_scanner_exception_test.cpp | 131 ++++++-
.../iceberg_delete_file_reader_helper_test.cpp | 16 +-
.../format/table/iceberg/iceberg_reader_test.cpp | 73 ++++
.../format_v2/orc/orc_file_input_stream_test.cpp | 33 ++
be/test/format_v2/orc/orc_reader_test.cpp | 34 ++
be/test/format_v2/table/iceberg_reader_test.cpp | 120 ++++++-
be/test/io/fs/file_handle_cache_test.cpp | 391 +++++++++++++++++++++
be/test/io/fs/hdfs_file_system_test.cpp | 70 ++++
.../index/snii/snii_blob_directory_test.cpp | 36 ++
.../segment/inverted_index_file_reader_test.cpp | 27 ++
.../connector/iceberg/IcebergScanPlanProvider.java | 6 +-
.../doris/connector/iceberg/IcebergScanRange.java | 23 +-
.../iceberg/IcebergScanPlanProviderTest.java | 108 ++++++
.../connector/iceberg/IcebergScanRangeTest.java | 72 +++-
gensrc/thrift/PlanNodes.thrift | 1 +
31 files changed, 1266 insertions(+), 80 deletions(-)
diff --git a/be/src/exec/scan/file_scanner.cpp
b/be/src/exec/scan/file_scanner.cpp
index b367a799459..d9306a37f76 100644
--- a/be/src/exec/scan/file_scanner.cpp
+++ b/be/src/exec/scan/file_scanner.cpp
@@ -580,8 +580,14 @@ Status FileScanner::_get_block_wrapped(RuntimeState*
state, Block* block, bool*
// Read next block.
// Some of column in block may not be filled (column not exist in
file)
- RETURN_IF_ERROR(
- _cur_reader->get_next_block(_src_block_ptr, &read_rows,
&_cur_reader_eof));
+ Status st = _cur_reader->get_next_block(_src_block_ptr,
&read_rows, &_cur_reader_eof);
+ // Lazy open may surface NOT_FOUND on the first read; skip as
above.
+ if (st.is<ErrorCode::NOT_FOUND>() &&
config::ignore_not_found_file_in_external_table) {
+ _cur_reader_eof = true;
+ COUNTER_UPDATE(_not_found_file_counter, 1);
+ continue;
+ }
+ RETURN_IF_ERROR(st);
}
// use read_rows instead of _src_block_ptr->rows(), because the first
column of _src_block_ptr
// may not be filled after calling `get_next_block()`, so
_src_block_ptr->rows() may return wrong result.
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 f8828e402b0..a2539161e64 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,12 @@ 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;
+ range.file_size = file_size;
return range;
}
@@ -276,7 +276,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.__isset.file_size ?
delete_file.file_size : -1);
if (options.fs_name != nullptr && !options.fs_name->empty()) {
delete_range.__set_fs_name(*options.fs_name);
}
@@ -348,7 +349,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.__isset.file_size ?
delete_file.file_size : -1);
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 adc0ef196f4..d43443ecd8e 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<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);
diff --git a/be/src/format/table/iceberg_reader_mixin.h
b/be/src/format/table/iceberg_reader_mixin.h
index 437b76d0e2d..a9d0925f10e 100644
--- a/be/src/format/table/iceberg_reader_mixin.h
+++ b/be/src/format/table/iceberg_reader_mixin.h
@@ -143,6 +143,10 @@ public:
return _position_delete_base(data_file_path, delete_files);
}
+ Status TEST_read_equality_delete_file(const TIcebergDeleteFileDesc&
delete_file) {
+ return _read_equality_delete_file(delete_file);
+ }
+
void TEST_set_column_name_to_block_index(
std::unordered_map<std::string, uint32_t>*
column_name_to_block_index) {
this->col_name_to_block_idx_ref() = column_name_to_block_index;
@@ -848,7 +852,7 @@ Status
IcebergReaderMixin<BaseReader>::_read_equality_delete_file(
delete_desc.path = delete_file.path;
delete_desc.start_offset = 0;
delete_desc.size = -1;
- delete_desc.file_size = -1;
+ delete_desc.file_size = delete_file.__isset.file_size ?
delete_file.file_size : -1;
std::unique_ptr<GenericReader> reader =
_create_equality_reader(delete_desc);
RETURN_IF_ERROR(reader->init_schema_reader());
@@ -1290,7 +1294,8 @@ Status
IcebergReaderMixin<BaseReader>::_position_delete_base(
delete_file_range.path = delete_file.path;
delete_file_range.start_offset = 0;
delete_file_range.size = -1;
- delete_file_range.file_size = -1;
+ delete_file_range.file_size =
+ delete_file.__isset.file_size ?
delete_file.file_size : -1;
create_status =
_read_position_delete_file(&delete_file_range,
position_delete.get());
if (!create_status) {
diff --git a/be/src/format_v2/orc/orc_reader.cpp
b/be/src/format_v2/orc/orc_reader.cpp
index 83cb929e18b..b59eafbadfb 100644
--- a/be/src/format_v2/orc/orc_reader.cpp
+++ b/be/src/format_v2/orc/orc_reader.cpp
@@ -945,8 +945,15 @@ Status OrcReader::init(RuntimeState* state) {
if (is_orc_stop(_io_ctx.get(), e)) {
return Status::EndOfFile("stop");
}
+ // invoker maybe just skip Status.NotFound and continue
+ // so we need distinguish between it and other kinds of errors
+ const std::string err_msg = e.what();
+ if (err_msg.find("No such file or directory") != std::string::npos ||
+ err_msg.find("NoSuchKey") != std::string::npos) {
+ return Status::NotFound(err_msg);
+ }
return Status::InternalError("Failed to open ORC file {}: {}",
_file_description->path,
- e.what());
+ err_msg);
}
return Status::OK();
}
diff --git a/be/src/format_v2/table/iceberg_reader.cpp
b/be/src/format_v2/table/iceberg_reader.cpp
index 413c5382090..dd952eea603 100644
--- a/be/src/format_v2/table/iceberg_reader.cpp
+++ b/be/src/format_v2/table/iceberg_reader.cpp
@@ -1239,7 +1239,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->__isset.file_size ?
deletion_vector->file_size : -1;
desc->format = DeleteFileDesc::Format::ICEBERG;
*has_delete_file = true;
return Status::OK();
@@ -1595,7 +1595,8 @@ 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.__isset.file_size ?
delete_file.file_size : -1);
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 41617ba1015..56ce889416a 100644
--- a/be/src/io/fs/file_handle_cache.cpp
+++ b/be/src/io/fs/file_handle_cache.cpp
@@ -26,7 +26,11 @@
#include <tuple>
#include "common/cast_set.h"
+#include "common/metrics/doris_metrics.h"
+#include "cpp/sync_point.h"
#include "io/fs/err_utils.h"
+#include "io/hdfs_util.h"
+#include "util/bvar_helper.h"
#include "util/hash_util.hpp"
#include "util/time.h"
namespace doris::io {
@@ -34,29 +38,30 @@ namespace doris::io {
HdfsFileHandle::~HdfsFileHandle() {
if (_hdfs_file != nullptr && _fs != nullptr) {
VLOG_FILE << "hdfsCloseFile() fid=" << _hdfs_file;
- hdfsCloseFile(_fs, _hdfs_file); // TODO: check return code
+ SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_close_latency);
+ SYNC_POINT_HOOK_RETURN_VALUE(hdfsCloseFile(_fs, _hdfs_file),
+ "HdfsFileHandle::close::hdfsCloseFile");
+ DorisMetrics::instance()->hdfs_file_open_reading->increment(-1);
}
_fs = nullptr;
_hdfs_file = nullptr;
}
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) {
- hdfsFileInfo* file_info = hdfsGetPathInfo(_fs, _fname.c_str());
+ SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_get_path_info_latency);
+ auto* file_info = SYNC_POINT_HOOK_RETURN_VALUE(hdfsGetPathInfo(_fs,
_fname.c_str()),
+
"HdfsFileHandle::init::hdfsGetPathInfo");
if (file_info == nullptr) {
- return Status::InternalError("failed to get file size of {}: {}",
_fname, hdfs_error());
+ std::string err_msg =
+ SYNC_POINT_HOOK_RETURN_VALUE(hdfs_error(),
"HdfsFileHandle::init::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 get file size of {}: {}",
_fname, err_msg);
}
_file_size = file_info->mSize;
hdfsFreeFileInfo(file_info, 1);
@@ -64,6 +69,31 @@ 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;
+ SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_open_latency);
+ _hdfs_file =
+ SYNC_POINT_HOOK_RETURN_VALUE(hdfsOpenFile(_fs, _fname.c_str(),
O_RDONLY, 0, 0, 0),
+
"HdfsFileHandle::ensure_open::hdfsOpenFile");
+ if (_hdfs_file != nullptr) {
+ _open_status = Status::OK();
+ DorisMetrics::instance()->hdfs_file_open_reading->increment(1);
+ DorisMetrics::instance()->hdfs_file_reader_total->increment(1);
+ } else {
+ // Capture error inside the opening thread (libhdfs last-error is
thread-local).
+ std::string _err_msg = SYNC_POINT_HOOK_RETURN_VALUE(
+ hdfs_error(), "HdfsFileHandle::ensure_open::hdfs_error");
+ if (_err_msg.find("No such file or directory") !=
std::string::npos) {
+ _open_status = Status::NotFound(_err_msg);
+ } else {
+ _open_status = Status::InternalError("failed to open {}: {}",
_fname, _err_msg);
+ }
+ }
+ });
+ return _open_status;
+}
+
CachedHdfsFileHandle::CachedHdfsFileHandle(const hdfsFS& fs, const
std::string& fname,
int64_t mtime)
: HdfsFileHandle(fs, fname, mtime) {}
@@ -98,8 +128,16 @@ void FileHandleCache::Accessor::destroy() {
FileHandleCache::Accessor::~Accessor() {
if (_cache_accessor.get()) {
+ auto* handle = get();
+ if (handle->file() == nullptr) {
+ // Not opened or open failed (call_once won't retry), destroy to
avoid cache pollution
+ destroy();
+ return;
+ }
#ifdef USE_HADOOP_HDFS
- if (hdfsUnbufferFile(get()->file()) != 0) {
+ int unbuffer_ret =
SYNC_POINT_HOOK_RETURN_VALUE(hdfsUnbufferFile(handle->file()),
+
"HdfsFileHandle::close::hdfsUnbufferFile");
+ if (unbuffer_ret != 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 ce3c708ba99..46e91035746 100644
--- a/be/src/io/fs/file_handle_cache.h
+++ b/be/src/io/fs/file_handle_cache.h
@@ -26,6 +26,7 @@
#include <list>
#include <map>
#include <memory>
+#include <mutex>
#include <string>
#include <utility>
@@ -49,9 +50,12 @@ public:
/// 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,9 @@ private:
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;
+ Status _open_status;
};
/// 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 c9b3f70580c..6b02c273723 100644
--- a/be/src/io/fs/hdfs_file_reader.cpp
+++ b/be/src/io/fs/hdfs_file_reader.cpp
@@ -25,7 +25,6 @@
#include "bvar/latency_recorder.h"
#include "bvar/reducer.h"
#include "common/compiler_util.h" // IWYU pragma: keep
-#include "common/metrics/doris_metrics.h"
#include "cpp/sync_point.h"
#include "io/fs/err_utils.h"
#include "io/hdfs_util.h"
@@ -34,6 +33,7 @@
#include "runtime/workload_management/io_throttle.h"
#include "runtime/workload_management/resource_context.h"
#include "service/backend_options.h"
+#include "util/bvar_helper.h"
namespace doris::io {
@@ -74,9 +74,6 @@ HdfsFileReader::HdfsFileReader(Path path, std::string
fs_name, FileHandleCache::
_accessor(std::move(accessor)),
_mtime(mtime) {
_handle = _accessor.get();
-
- DorisMetrics::instance()->hdfs_file_open_reading->increment(1);
- DorisMetrics::instance()->hdfs_file_reader_total->increment(1);
}
HdfsFileReader::~HdfsFileReader() {
@@ -84,15 +81,20 @@ HdfsFileReader::~HdfsFileReader() {
}
Status HdfsFileReader::close() {
- bool expected = false;
- if (_closed.compare_exchange_strong(expected, true,
std::memory_order_acq_rel)) {
- DorisMetrics::instance()->hdfs_file_open_reading->increment(-1);
- }
+ _closed = true;
return Status::OK();
}
Status HdfsFileReader::read_at_impl(size_t offset, Slice result, size_t*
bytes_read,
const IOContext* io_ctx) {
+ if (closed()) [[unlikely]] {
+ return Status::InternalError("read closed file: {}", _path.native());
+ }
+ if (_handle == nullptr) [[unlikely]] {
+ return Status::InternalError("cached hdfs file handle has been
destroyed: {}",
+ _path.native());
+ }
+ RETURN_IF_ERROR(_handle->ensure_open());
auto st = do_read_at_impl(offset, result, bytes_read, io_ctx);
if (!st.ok()) {
_handle = nullptr;
@@ -104,15 +106,6 @@ Status HdfsFileReader::read_at_impl(size_t offset, Slice
result, size_t* bytes_r
#ifdef USE_HADOOP_HDFS
Status HdfsFileReader::do_read_at_impl(size_t offset, Slice result, size_t*
bytes_read,
const IOContext* /*io_ctx*/) {
- if (closed()) [[unlikely]] {
- return Status::InternalError("read closed file: {}", _path.native());
- }
-
- if (_handle == nullptr) [[unlikely]] {
- return Status::InternalError("cached hdfs file handle has been
destroyed: {}",
- _path.native());
- }
-
if (offset > _handle->file_size()) {
return Status::IOError("offset exceeds file size(offset: {}, file
size: {}, path: {})",
offset, _handle->file_size(), _path.native());
@@ -133,8 +126,12 @@ Status HdfsFileReader::do_read_at_impl(size_t offset,
Slice result, size_t* byte
int64_t max_to_read = bytes_req - has_read;
tSize to_read = static_cast<tSize>(
std::min(max_to_read,
static_cast<int64_t>(std::numeric_limits<tSize>::max())));
- tSize loop_read = hdfsPread(_handle->fs(), _handle->file(), offset +
has_read,
- to + has_read, to_read);
+ tSize loop_read;
+ {
+ SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_read_latency);
+ loop_read = hdfsPread(_handle->fs(), _handle->file(), offset +
has_read, to + has_read,
+ to_read);
+ }
{
[[maybe_unused]] Status error_ret;
TEST_INJECTION_POINT_RETURN_WITH_VALUE("HdfsFileReader:read_error", error_ret);
@@ -166,10 +163,6 @@ Status HdfsFileReader::do_read_at_impl(size_t offset,
Slice result, size_t* byte
// TODO: rethink here to see if there are some difference between hdfsPread()
and hdfsRead()
Status HdfsFileReader::do_read_at_impl(size_t offset, Slice result, size_t*
bytes_read,
const IOContext* /*io_ctx*/) {
- if (closed()) [[unlikely]] {
- return Status::InternalError("read closed file: ", _path.native());
- }
-
if (offset > _handle->file_size()) {
return Status::IOError("offset exceeds file size(offset: {}, file
size: {}, path: {})",
offset, _handle->file_size(), _path.native());
@@ -199,8 +192,12 @@ Status HdfsFileReader::do_read_at_impl(size_t offset,
Slice result, size_t* byte
size_t has_read = 0;
while (has_read < bytes_req) {
- int64_t loop_read = hdfsRead(_handle->fs(), _handle->file(), to +
has_read,
- static_cast<int32_t>(bytes_req -
has_read));
+ int64_t loop_read;
+ {
+ SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_read_latency);
+ loop_read = hdfsRead(_handle->fs(), _handle->file(), to + has_read,
+ static_cast<int32_t>(bytes_req - has_read));
+ }
if (loop_read < 0) {
// invoker maybe just skip Status.NotFound and continue
// so we need distinguish between it and other kinds of errors
diff --git a/be/src/io/fs/hdfs_file_system.cpp
b/be/src/io/fs/hdfs_file_system.cpp
index e85b582e24f..5043c9eaaa6 100644
--- a/be/src/io/fs/hdfs_file_system.cpp
+++ b/be/src/io/fs/hdfs_file_system.cpp
@@ -40,6 +40,7 @@
#include "io/hdfs_builder.h"
#include "io/hdfs_util.h"
#include "runtime/exec_env.h"
+#include "util/bvar_helper.h"
#include "util/obj_lru_cache.h"
#include "util/slice.h"
@@ -117,7 +118,11 @@ Status HdfsFileSystem::open_file_internal(const Path&
file, FileReaderSPtr* read
Status HdfsFileSystem::create_directory_impl(const Path& dir, bool
failed_if_exists) {
CHECK_HDFS_HANDLER(_fs_handler);
Path real_path = convert_path(dir, _fs_name);
- int res = hdfsCreateDirectory(_fs_handler->hdfs_fs,
real_path.string().c_str());
+ int res;
+ {
+ SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_create_dir_latency);
+ res = hdfsCreateDirectory(_fs_handler->hdfs_fs,
real_path.string().c_str());
+ }
if (res == -1) {
return Status::IOError("failed to create directory {}: {}",
dir.native(), hdfs_error());
}
@@ -180,6 +185,7 @@ Status HdfsFileSystem::exists_impl(const Path& path, bool*
res) const {
Status HdfsFileSystem::file_size_impl(const Path& path, int64_t* file_size)
const {
CHECK_HDFS_HANDLER(_fs_handler);
Path real_path = convert_path(path, _fs_name);
+ SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_get_path_info_latency);
hdfsFileInfo* file_info = hdfsGetPathInfo(_fs_handler->hdfs_fs,
real_path.string().c_str());
if (file_info == nullptr) {
return Status::IOError("failed to get file size of {}: {}",
path.native(), hdfs_error());
@@ -222,6 +228,7 @@ Status HdfsFileSystem::list_impl(const Path& path, bool
only_file, std::vector<F
}
Status HdfsFileSystem::rename_impl(const Path& orig_name, const Path&
new_name) {
+ CHECK_HDFS_HANDLER(_fs_handler);
Path normal_orig_name = convert_path(orig_name, _fs_name);
Path normal_new_name = convert_path(new_name, _fs_name);
int ret = hdfsRename(_fs_handler->hdfs_fs, normal_orig_name.c_str(),
normal_new_name.c_str());
diff --git a/be/src/io/fs/local_file_reader.cpp
b/be/src/io/fs/local_file_reader.cpp
index 50bac007224..4f7628af166 100644
--- a/be/src/io/fs/local_file_reader.cpp
+++ b/be/src/io/fs/local_file_reader.cpp
@@ -176,7 +176,7 @@ Status LocalFileReader::read_at_impl(size_t offset, Slice
result, size_t* bytes_
if ((sub_path.empty() && _path.filename().compare(kTestFilePath))
||
(!sub_path.empty() && _path.native().find(sub_path) !=
std::string::npos)) {
res = -1;
- errno = EIO;
+ errno = dp->param<int>("errno", EIO);
LOG(WARNING) << Status::IOError("debug read io error: {}",
_path.native());
}
});
diff --git a/be/src/io/hdfs_util.cpp b/be/src/io/hdfs_util.cpp
index c7060108282..fb54f82b6a1 100644
--- a/be/src/io/hdfs_util.cpp
+++ b/be/src/io/hdfs_util.cpp
@@ -41,6 +41,7 @@ bvar::LatencyRecorder hdfs_close_latency("hdfs_close");
bvar::LatencyRecorder hdfs_flush_latency("hdfs_flush");
bvar::LatencyRecorder hdfs_hflush_latency("hdfs_hflush");
bvar::LatencyRecorder hdfs_hsync_latency("hdfs_hsync");
+bvar::LatencyRecorder hdfs_get_path_info_latency("hdfs_get_path_info");
}; // namespace hdfs_bvar
Path convert_path(const Path& path, const std::string& namenode) {
diff --git a/be/src/io/hdfs_util.h b/be/src/io/hdfs_util.h
index 8d63a19ec92..d7a74134db6 100644
--- a/be/src/io/hdfs_util.h
+++ b/be/src/io/hdfs_util.h
@@ -80,6 +80,7 @@ extern bvar::LatencyRecorder hdfs_close_latency;
extern bvar::LatencyRecorder hdfs_flush_latency;
extern bvar::LatencyRecorder hdfs_hflush_latency;
extern bvar::LatencyRecorder hdfs_hsync_latency;
+extern bvar::LatencyRecorder hdfs_get_path_info_latency;
}; // namespace hdfs_bvar
// if the format of path is hdfs://ip:port/path, replace it to /path.
diff --git a/be/src/storage/index/index_file_reader.cpp
b/be/src/storage/index/index_file_reader.cpp
index ac2f17db252..4b51caf6ace 100644
--- a/be/src/storage/index/index_file_reader.cpp
+++ b/be/src/storage/index/index_file_reader.cpp
@@ -133,6 +133,11 @@ Status IndexFileReader::_init_from(int32_t
read_buffer_size, const io::IOContext
index_file_full_path, err.what());
}
}
+ // Lazy open can surface a missing file as a read error; keep NotFound
distinguishable
+ if (err.number() == CL_ERR_FileNotFound) {
+ return Status::Error<ErrorCode::INVERTED_INDEX_FILE_NOT_FOUND>(
+ "inverted index file {} is not found.",
index_file_full_path);
+ }
return Status::Error<ErrorCode::INVERTED_INDEX_CLUCENE_ERROR>(
"CLuceneError occur when init idx file {}, error msg: {}",
index_file_full_path,
err.what());
@@ -318,6 +323,11 @@ Result<std::unique_ptr<DorisCompoundReader,
DirectoryDeleter>> IndexFileReader::
// 3. read file in DorisCompoundReader
compound_reader.reset(new DorisCompoundReader(index_input,
_read_buffer_size));
} catch (CLuceneError& err) {
+ // Lazy open can surface a missing file as a read error; keep
NotFound distinguishable
+ if (err.number() == CL_ERR_FileNotFound) {
+ return
ResultError(Status::Error<ErrorCode::INVERTED_INDEX_FILE_NOT_FOUND>(
+ "inverted index file {} is not found.",
index_file_path));
+ }
return
ResultError(Status::Error<ErrorCode::INVERTED_INDEX_CLUCENE_ERROR>(
"CLuceneError occur when open idx file {}, error msg: {}",
index_file_path,
err.what()));
diff --git a/be/src/storage/index/inverted/inverted_index_fs_directory.cpp
b/be/src/storage/index/inverted/inverted_index_fs_directory.cpp
index dcb20b73445..37e2595da71 100644
--- a/be/src/storage/index/inverted/inverted_index_fs_directory.cpp
+++ b/be/src/storage/index/inverted/inverted_index_fs_directory.cpp
@@ -241,7 +241,11 @@ void DorisFSDirectory::FSIndexInput::readInternal(uint8_t*
b, const int32_t len)
"DorisFSDirectory::FSIndexInput::readInternal_reader_read_at_error");
})
if (!st.ok()) {
- _CLTHROWA(CL_ERR_IO, "read past EOF");
+ // Carry NotFound across the CLucene boundary so index callers can
downgrade
+ if (st.is<ErrorCode::NOT_FOUND>()) {
+ _CLTHROWA(CL_ERR_FileNotFound,
st.to_string_no_stack().c_str());
+ }
+ _CLTHROWA(CL_ERR_IO, st.to_string_no_stack().c_str());
}
bufferLength = len;
DBUG_EXECUTE_IF("DorisFSDirectory::FSIndexInput::readInternal_bytes_read_error",
diff --git a/be/src/storage/index/snii/snii_blob_directory.cpp
b/be/src/storage/index/snii/snii_blob_directory.cpp
index 628ab0ef85e..1c1b0db1c7f 100644
--- a/be/src/storage/index/snii/snii_blob_directory.cpp
+++ b/be/src/storage/index/snii/snii_blob_directory.cpp
@@ -84,8 +84,12 @@ protected:
static_cast<size_t>(len));
}
if (!status.ok()) {
+ // Lazy open can surface a missing file as a read error; keep
NotFound distinguishable
+ if (status.is<ErrorCode::NOT_FOUND>()) {
+ _CLTHROWA(CL_ERR_FileNotFound,
status.to_string_no_stack().c_str());
+ }
// CLuceneError STRDUPs the message; the temporary is safe.
- _CLTHROWA(CL_ERR_IO, status.to_string().c_str());
+ _CLTHROWA(CL_ERR_IO, status.to_string_no_stack().c_str());
}
}
void seekInternal(const int64_t /*pos*/) override {}
diff --git a/be/test/exec/scan/vfile_scanner_exception_test.cpp
b/be/test/exec/scan/vfile_scanner_exception_test.cpp
index 70b3d07f8ef..62bc3d73011 100644
--- a/be/test/exec/scan/vfile_scanner_exception_test.cpp
+++ b/be/test/exec/scan/vfile_scanner_exception_test.cpp
@@ -23,6 +23,7 @@
#include <utility>
#include <vector>
+#include "common/config.h"
#include "common/object_pool.h"
#include "core/data_type/data_type_number.h"
#include "core/data_type/data_type_string.h"
@@ -31,13 +32,17 @@
#include "exec/scan/file_scanner.h"
#include "exec/scan/split_source_connector.h"
#include "format_v2/table/hive_reader.h"
+#include "io/fs/hdfs/hdfs_mgr.h"
#include "io/fs/local_file_system.h"
+#include "io/hdfs_util.h"
#include "load/group_commit/wal/wal_manager.h"
#include "runtime/cluster_info.h"
#include "runtime/descriptors.h"
+#include "runtime/exec_env.h"
#include "runtime/memory/mem_tracker.h"
#include "runtime/runtime_state.h"
#include "runtime/user_function_cache.h"
+#include "util/defer_op.h"
namespace doris {
class TestSplitSourceConnectorStub : public SplitSourceConnector {
@@ -65,6 +70,47 @@ public:
TFileScanRangeParams* get_params() override { return &_scan_range.params; }
};
+// Returns a fake HDFS handler so reader construction never touches the
network.
+class FakeHdfsMgr final : public io::HdfsMgr {
+public:
+ Status _create_hdfs_fs_impl(const THdfsParams& hdfs_params, const
std::string& fs_name,
+ std::shared_ptr<io::HdfsHandler>* fs_handler)
override {
+ *fs_handler = std::make_shared<io::HdfsHandler>(
+ reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1)), false,
"", "", fs_name);
+ return Status::OK();
+ }
+};
+
+// Registers a SyncPoint callback that returns a fixed value.
+template <typename T>
+static void set_mock_return(const std::string& point, T value,
SyncPoint::CallbackGuard* guard) {
+ SyncPoint::get_instance()->set_call_back(
+ point,
+ [value = std::move(value)](auto&& args) {
+ auto* ret = try_any_cast_ret<T>(args);
+ ret->first = std::move(value);
+ ret->second = true;
+ },
+ guard);
+}
+
+// Injects a read-stage NOT_FOUND through the HDFS lazy-open sync points.
+struct HdfsNotFoundGuard {
+ SyncPoint::CallbackGuard open_guard;
+ SyncPoint::CallbackGuard err_guard;
+ SyncPoint::CallbackGuard close_guard;
+ HdfsNotFoundGuard() {
+ auto* sp = SyncPoint::get_instance();
+ sp->enable_processing();
+ set_mock_return<hdfsFile>("HdfsFileHandle::ensure_open::hdfsOpenFile",
nullptr,
+ &open_guard);
+ set_mock_return<std::string>("HdfsFileHandle::ensure_open::hdfs_error",
+ "No such file or directory", &err_guard);
+ set_mock_return<int>("HdfsFileHandle::close::hdfsCloseFile", 0,
&close_guard);
+ }
+ ~HdfsNotFoundGuard() { SyncPoint::get_instance()->disable_processing(); }
+};
+
class VfileScannerExceptionTest : public testing::Test {
public:
VfileScannerExceptionTest()
@@ -78,17 +124,26 @@ public:
}
void init();
void generate_scanner(std::shared_ptr<FileScanner>& scanner);
+ // Fake HDFS CSV range with known size, so open is deferred to the first
read.
+ void prepare_hdfs_csv_range();
void TearDown() override {
WARN_IF_ERROR(_scan_node->close(&_runtime_state), "fail to close
scan_node")
+ // Avoid leaving a dangling _hdfs_mgr after _fake_hdfs_mgr is
destroyed.
+ ExecEnv::GetInstance()->_hdfs_mgr = _old_hdfs_mgr;
}
protected:
- virtual void SetUp() override {}
+ void SetUp() override {
+ _old_hdfs_mgr = ExecEnv::GetInstance()->_hdfs_mgr;
+ ExecEnv::GetInstance()->_hdfs_mgr = &_fake_hdfs_mgr;
+ }
private:
void _init_desc_table();
+ FakeHdfsMgr _fake_hdfs_mgr;
+ io::HdfsMgr* _old_hdfs_mgr = nullptr;
ExecEnv* _env = nullptr;
int64_t _backend_id = 1001;
std::string _label_1 = "test1";
@@ -286,6 +341,34 @@ void
VfileScannerExceptionTest::generate_scanner(std::shared_ptr<FileScanner>& s
WARN_IF_ERROR(scanner->init(&_runtime_state, _conjuncts), "fail to prepare
scanner");
}
+void VfileScannerExceptionTest::prepare_hdfs_csv_range() {
+ _range_desc.path = "hdfs://fake-nn:8020/not_found/data.csv";
+ _range_desc.start_offset = 0;
+ _range_desc.size = 100;
+ _range_desc.__set_file_size(100);
+ _ranges[0] = _range_desc;
+ _scan_range.ranges = _ranges;
+ auto& params = _scan_range.params;
+ params.format_type = TFileFormatType::FORMAT_CSV_PLAIN;
+ params.file_type = TFileType::FILE_HDFS;
+ params.hdfs_params.__set_fs_name("hdfs://fake-nn:8020");
+ params.__isset.file_attributes = true;
+ params.file_attributes.__isset.text_params = true;
+ params.file_attributes.text_params.column_separator = ",";
+ params.file_attributes.text_params.line_delimiter = "\n";
+ params.__isset.column_idxs = true;
+ params.column_idxs = {0, 1, 2};
+ params.__set_num_of_columns_from_file(3);
+ params.__isset.required_slots = true;
+ params.required_slots.clear();
+ for (int32_t slot_id = 1; slot_id <= 3; ++slot_id) {
+ TFileScanSlotInfo slot_info;
+ slot_info.__set_slot_id(slot_id);
+ slot_info.__set_is_file_slot(true);
+ params.required_slots.push_back(slot_info);
+ }
+}
+
TEST_F(VfileScannerExceptionTest, failure_case) {
std::shared_ptr<FileScanner> scanner = nullptr;
generate_scanner(scanner);
@@ -341,6 +424,52 @@ TEST_F(VfileScannerExceptionTest,
process_late_arrival_conjuncts_retain) {
WARN_IF_ERROR(scanner->close(&_runtime_state), "fail to close scanner");
}
+// A lazy-open NOT_FOUND on the first HDFS read must be skipped and counted,
not fail the scan.
+TEST_F(VfileScannerExceptionTest, read_stage_not_found_skipped_when_enabled) {
+ HdfsNotFoundGuard guard;
+ const bool old_ignore = config::ignore_not_found_file_in_external_table;
+ config::ignore_not_found_file_in_external_table = true;
+ Defer restore_ignore {[&]() {
config::ignore_not_found_file_in_external_table = old_ignore; }};
+
+ prepare_hdfs_csv_range();
+ std::shared_ptr<FileScanner> scanner = nullptr;
+ generate_scanner(scanner);
+
+ std::unique_ptr<Block> block(new Block());
+ bool eof = false;
+ auto st = scanner->get_block(&_runtime_state, block.get(), &eof);
+ EXPECT_TRUE(st.ok()) << st;
+ EXPECT_TRUE(eof);
+ EXPECT_EQ(block->rows(), 0);
+
+ auto* local_state =
&(_runtime_state.get_local_state(0)->cast<FileScanLocalState>());
+ auto* counter =
local_state->scanner_profile()->get_counter("NotFoundFileNum");
+ ASSERT_NE(counter, nullptr);
+ EXPECT_EQ(counter->value(), 1);
+
+ WARN_IF_ERROR(scanner->close(&_runtime_state), "fail to close scanner");
+}
+
+// With the skip config off, a lazy-open NOT_FOUND must surface as an error,
not be wrapped.
+TEST_F(VfileScannerExceptionTest, read_stage_not_found_errors_when_disabled) {
+ HdfsNotFoundGuard guard;
+ const bool old_ignore = config::ignore_not_found_file_in_external_table;
+ config::ignore_not_found_file_in_external_table = false;
+ Defer restore_ignore {[&]() {
config::ignore_not_found_file_in_external_table = old_ignore; }};
+
+ prepare_hdfs_csv_range();
+ std::shared_ptr<FileScanner> scanner = nullptr;
+ generate_scanner(scanner);
+
+ std::unique_ptr<Block> block(new Block());
+ bool eof = false;
+ auto st = scanner->get_block(&_runtime_state, block.get(), &eof);
+ EXPECT_FALSE(st.ok());
+ EXPECT_TRUE(st.is<ErrorCode::NOT_FOUND>()) << st;
+
+ WARN_IF_ERROR(scanner->close(&_runtime_state), "fail to close scanner");
+}
+
TEST(HiveReaderPositionMappingTest, PositionMappingUsesColumnIdxsForFileSlots)
{
TQueryOptions query_options;
query_options.hive_parquet_use_column_names = false;
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 841e01f9cc9..72b519e2191 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
@@ -214,11 +214,14 @@ 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);
EXPECT_EQ(range.file_size, -1);
+
+ auto range2 =
build_iceberg_delete_file_range("s3://bucket/delete.parquet", 1024);
+ EXPECT_EQ(range2.file_size, 1024);
}
TEST(IcebergDeleteFileReaderHelperTest, IsDeletionVector) {
@@ -384,7 +387,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,
@@ -392,6 +395,15 @@ TEST(IcebergDeleteFileReaderHelperTest,
DeletionVectorReaderValidatesOpenedFileR
const auto exact_status = exact_reader.open();
EXPECT_TRUE(exact_status.ok()) << exact_status;
+ // Manifest-style range: exact file_size provided by the FE instead of
-1.
+ TFileRangeDesc manifest_range =
build_iceberg_delete_file_range(dv_path, dv_size);
+ manifest_range.start_offset = 4;
+ manifest_range.size = dv_size - manifest_range.start_offset;
+ DeletionVectorReader manifest_reader(&state, &profile, scan_params,
manifest_range,
+ &io_context.io_ctx);
+ const auto manifest_status = manifest_reader.open();
+ EXPECT_TRUE(manifest_status.ok()) << manifest_status;
+
TFileRangeDesc oversized_range = exact_range;
oversized_range.size = MAX_ICEBERG_DELETION_VECTOR_BYTES;
DeletionVectorReader oversized_reader(&state, &profile, scan_params,
oversized_range,
diff --git a/be/test/format/table/iceberg/iceberg_reader_test.cpp
b/be/test/format/table/iceberg/iceberg_reader_test.cpp
index ef1c469ca23..0dee309fb7f 100644
--- a/be/test/format/table/iceberg/iceberg_reader_test.cpp
+++ b/be/test/format/table/iceberg/iceberg_reader_test.cpp
@@ -1718,6 +1718,79 @@ TEST_F(IcebergReaderTest,
v1_position_delete_read_error_releases_cache_entry) {
EXPECT_NE(status.to_string().find(delete_file.path), std::string::npos);
}
+// An inflated size breaks footer reads, proving the v1 position-delete path
consumes the FE file_size.
+TEST_F(IcebergReaderTest, v1_position_delete_consumes_delete_file_size) {
+ 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<ParquetReader>::POSITION_DELETE);
+ delete_file.__set_path(mixed_position_delete_file());
+ delete_file.__set_file_size(1 << 30); // wrong on purpose: must reach the
reader
+
+ const auto status =
+
iceberg_reader.TEST_position_delete_base("file:///tmp/data.parquet",
{delete_file});
+
+ ASSERT_FALSE(status.ok());
+}
+
+// An inflated size must break the read, proving the v1 equality-delete path
consumes the FE file_size.
+TEST_F(IcebergReaderTest, v1_equality_delete_consumes_delete_file_size) {
+ const auto test_dir = std::filesystem::temp_directory_path() /
"doris_v1_eq_delete_size_test";
+ std::filesystem::remove_all(test_dir);
+ std::filesystem::create_directories(test_dir);
+ const auto delete_file_path = (test_dir /
"equality-delete.parquet").string();
+ write_iceberg_int_equality_delete_parquet_file(delete_file_path, "id", 0,
2);
+
+ 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<ParquetReader>::EQUALITY_DELETE);
+ delete_file.__set_path(delete_file_path);
+ delete_file.__set_field_ids({0});
+ delete_file.__set_file_format(TFileFormatType::FORMAT_PARQUET);
+ delete_file.__set_file_size(1 << 30); // wrong on purpose: must reach the
reader
+
+ const auto status =
iceberg_reader.TEST_read_equality_delete_file(delete_file);
+
+ ASSERT_FALSE(status.ok());
+}
+
TEST_F(IcebergReaderTest,
v1_rejects_missing_required_field_without_initial_default) {
schema::external::TField field;
field.__set_name("required_added");
diff --git a/be/test/format_v2/orc/orc_file_input_stream_test.cpp
b/be/test/format_v2/orc/orc_file_input_stream_test.cpp
index 8d63a4ab21a..bfc56f0349f 100644
--- a/be/test/format_v2/orc/orc_file_input_stream_test.cpp
+++ b/be/test/format_v2/orc/orc_file_input_stream_test.cpp
@@ -75,6 +75,25 @@ private:
io::Path _path = "/tmp/orc_v2_input_stream";
};
+// FileReader whose reads always fail with NotFound, mimicking a missing HDFS
file under lazy open.
+class NotFoundFileReader final : public io::FileReader {
+public:
+ Status close() override { return Status::OK(); }
+ const io::Path& path() const override { return _path; }
+ size_t size() const override { return 4096; }
+ bool closed() const override { return false; }
+ int64_t mtime() const override { return 0; }
+
+protected:
+ Status read_at_impl(size_t offset, Slice result, size_t* bytes_read,
+ const io::IOContext* io_ctx) override {
+ return Status::NotFound("failed to open /test/missing.orc: No such
file or directory");
+ }
+
+private:
+ io::Path _path = "/test/missing.orc";
+};
+
class TestStreamInformation final : public ::orc::StreamInformation {
public:
TestStreamInformation(TestStream stream, uint64_t offset) :
_stream(stream), _offset(offset) {}
@@ -374,5 +393,19 @@ TEST(OrcFileInputStreamTest,
MergedIoChildrenStayIsolatedInBothInitializationOrd
}
}
+// A failed read must carry the status text so the ORC reader can restore
NotFound (parity with v1).
+TEST(OrcFileInputStreamTest, ReadFailureCarriesNotFoundText) {
+ auto reader = std::make_shared<NotFoundFileReader>();
+ OrcFileInputStream input("missing.orc", reader, nullptr, nullptr, {});
+ std::array<char, 16> buf {};
+ try {
+ input.read(buf.data(), buf.size(), 0);
+ FAIL() << "expected orc::ParseError";
+ } catch (const ::orc::ParseError& e) {
+ const std::string msg = e.what();
+ EXPECT_NE(msg.find("No such file or directory"), std::string::npos);
+ }
+}
+
} // namespace
} // namespace doris::format::orc
diff --git a/be/test/format_v2/orc/orc_reader_test.cpp
b/be/test/format_v2/orc/orc_reader_test.cpp
index 1058d774d44..67b5d2a1676 100644
--- a/be/test/format_v2/orc/orc_reader_test.cpp
+++ b/be/test/format_v2/orc/orc_reader_test.cpp
@@ -18,6 +18,7 @@
#include "format_v2/orc/orc_reader.h"
#include <cctz/time_zone.h>
+#include <errno.h>
#include <gtest/gtest.h>
#include <unistd.h>
@@ -4732,6 +4733,39 @@ TEST_F(NewOrcReaderTest,
AggregatePushdownReturnsCountFromFileMetadata) {
EXPECT_TRUE(aggregate_result.columns.empty());
}
+// Only ENOENT-style errors map to NotFound so FileScannerV2 does not silently
skip unhealthy splits.
+TEST_F(NewOrcReaderTest, InitKeepsInternalErrorForDirectory) {
+ auto system_properties = std::make_shared<io::FileSystemProperties>();
+ system_properties->system_type = TFileType::FILE_LOCAL;
+ auto file_description = std::make_unique<io::FileDescription>();
+ file_description->path = _test_dir; // open() on a directory succeeds,
read fails with EISDIR
+ file_description->file_size = 4096;
+ format::orc::OrcReader reader(system_properties, file_description,
nullptr, nullptr,
+ std::nullopt);
+ RuntimeState state {TQueryOptions(), TQueryGlobals()};
+ auto st = reader.init(&state);
+ ASSERT_FALSE(st.ok());
+ EXPECT_TRUE(st.is<ErrorCode::INTERNAL_ERROR>()) << st;
+}
+
+// ENOENT injected at the file layer surfaces as NotFound so FileScannerV2 can
skip the split.
+TEST_F(NewOrcReaderTest, InitRestoresNotFoundFromReadFailure) {
+ const auto old_enable = config::enable_debug_points;
+ config::enable_debug_points = true;
+ const std::string point = "LocalFileReader::read_at_impl.io_error";
+ DebugPoints::instance()->add_with_params(point, {{"errno",
std::to_string(ENOENT)}});
+
+ auto reader = create_reader();
+ RuntimeState state {TQueryOptions(), TQueryGlobals()};
+ auto st = reader->init(&state);
+
+ DebugPoints::instance()->remove(point);
+ config::enable_debug_points = old_enable;
+
+ ASSERT_FALSE(st.ok());
+ EXPECT_TRUE(st.is<ErrorCode::NOT_FOUND>()) << st;
+}
+
TEST_F(NewOrcReaderTest, AggregatePushdownCountUsesOnlySplitStripes) {
const auto multi_stripe_file_path = (_test_dir /
"aggregate_count_split.orc").string();
write_multi_stripe_orc_int_file(multi_stripe_file_path);
diff --git a/be/test/format_v2/table/iceberg_reader_test.cpp
b/be/test/format_v2/table/iceberg_reader_test.cpp
index 4bbfe6dd992..f91c237f575 100644
--- a/be/test/format_v2/table/iceberg_reader_test.cpp
+++ b/be/test/format_v2/table/iceberg_reader_test.cpp
@@ -1290,12 +1290,17 @@ TIcebergDeleteFileDesc
make_iceberg_position_delete_file(const std::string& path
TIcebergDeleteFileDesc make_iceberg_equality_delete_file(
const std::string& path, const std::vector<int32_t>& field_ids,
- TFileFormatType::type file_format = TFileFormatType::FORMAT_PARQUET) {
+ TFileFormatType::type file_format = TFileFormatType::FORMAT_PARQUET,
+ int64_t file_size = -1) {
TIcebergDeleteFileDesc delete_file;
delete_file.__set_content(2);
delete_file.__set_path(path);
delete_file.__set_field_ids(field_ids);
delete_file.__set_file_format(file_format);
+ // Default callers emulate an old FE that sends no file_size.
+ if (file_size >= 0) {
+ delete_file.__set_file_size(file_size);
+ }
return delete_file;
}
@@ -5416,5 +5421,118 @@ TEST(IcebergV2ReaderTest,
DataFileIsMarkedImmutableForPageCache) {
EXPECT_TRUE(reader.current_data_file_is_immutable());
}
+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>()));
+
+ 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());
+
+ // A truncated size must break the read; a correct size is
indistinguishable from the stat fallback.
+ 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},
+
TFileFormatType::FORMAT_PARQUET,
+ delete_file_size -
16)}));
+
+ const bool prepare_ok = reader.prepare_split(split_options).ok();
+ Block block = build_table_block(projected_columns);
+ bool eos = false;
+ const bool read_ok = prepare_ok && reader.get_block(&block, &eos).ok();
+ ASSERT_FALSE(read_ok) << "truncated file_size must fail the delete-file
read";
+
+ ASSERT_TRUE(reader.close().ok());
+ std::filesystem::remove_all(test_dir);
+}
+
+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);
+
+ 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);
+
+ std::vector<ColumnDefinition> projected_columns;
+ projected_columns.push_back(make_table_column(0, "id",
std::make_shared<DataTypeInt32>()));
+
+ 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, {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
diff --git a/be/test/io/fs/file_handle_cache_test.cpp
b/be/test/io/fs/file_handle_cache_test.cpp
index 5c1f7d1d9e0..a7da38ad91a 100644
--- a/be/test/io/fs/file_handle_cache_test.cpp
+++ b/be/test/io/fs/file_handle_cache_test.cpp
@@ -19,8 +19,19 @@
#include <gtest/gtest.h>
+#include <atomic>
#include <cstdint>
+#include <cstdlib>
+#include <cstring>
#include <string>
+#include <thread>
+#include <utility>
+#include <vector>
+
+#include "cpp/sync_point.h"
+#include "gen_cpp/Status_types.h"
+#include "io/fs/hdfs_file_reader.h"
+#include "util/defer_op.h"
namespace doris::io {
@@ -40,4 +51,384 @@ TEST(FileHandleCacheTest, CacheKeyIncludesHdfsFs) {
mtime + 1));
}
+// init(file_size>0) does not open the file.
+TEST(FileHandleCacheTest, InitWithKnownFileSizeDoesNotOpenFile) {
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ ExclusiveHdfsFileHandle handle(mock_fs, "/nonexistent/file.parquet",
12345);
+ auto st = handle.init(4096);
+ ASSERT_TRUE(st.ok()) << st;
+ EXPECT_EQ(handle.file_size(), 4096);
+ EXPECT_EQ(handle.file(), nullptr);
+}
+
+// Destructor is safe when file was never opened.
+TEST(FileHandleCacheTest, DestructorSafeWithoutOpen) {
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ {
+ ExclusiveHdfsFileHandle handle(mock_fs, "/nonexistent/file.parquet",
12345);
+ ASSERT_TRUE(handle.init(4096).ok());
+ EXPECT_EQ(handle.file(), nullptr);
+ }
+}
+
+// Register a SyncPoint callback that returns a fixed value.
+template <typename T>
+static void set_mock_return(const std::string& point, T value,
SyncPoint::CallbackGuard* guard) {
+ SyncPoint::get_instance()->set_call_back(
+ point,
+ [value = std::move(value)](auto&& args) {
+ auto* ret = try_any_cast_ret<T>(args);
+ ret->first = std::move(value);
+ ret->second = true;
+ },
+ guard);
+}
+
+// Mocks hdfsOpenFile/hdfsCloseFile/hdfsGetPathInfo/hdfsUnbufferFile via
SyncPoint to avoid JNI.
+struct MockHandleGuard {
+ SyncPoint::CallbackGuard open_guard;
+ SyncPoint::CallbackGuard close_guard;
+ SyncPoint::CallbackGuard info_guard;
+ SyncPoint::CallbackGuard unbuffer_guard;
+ MockHandleGuard(hdfsFile mock_file, int64_t file_size = 4096) {
+ auto* sp = SyncPoint::get_instance();
+ sp->enable_processing();
+ set_mock_return<hdfsFile>("HdfsFileHandle::ensure_open::hdfsOpenFile",
mock_file,
+ &open_guard);
+ set_mock_return<int>("HdfsFileHandle::close::hdfsCloseFile", 0,
&close_guard);
+ // Each hit allocates a fresh heap info so init()'s real
hdfsFreeFileInfo can free it safely.
+ sp->set_call_back(
+ "HdfsFileHandle::init::hdfsGetPathInfo",
+ [file_size](auto&& args) {
+ auto* ret = try_any_cast_ret<hdfsFileInfo*>(args);
+ auto* info = static_cast<hdfsFileInfo*>(calloc(1,
sizeof(hdfsFileInfo)));
+ info->mSize = file_size;
+ info->mName = strdup("/test/mock-file.parquet");
+ ret->first = info;
+ ret->second = true;
+ },
+ &info_guard);
+ set_mock_return<int>("HdfsFileHandle::close::hdfsUnbufferFile", 0,
&unbuffer_guard);
+ }
+ ~MockHandleGuard() { SyncPoint::get_instance()->disable_processing(); }
+};
+
+// ensure_open() succeeds via SyncPoint mock.
+TEST(FileHandleCacheTest, EnsureOpenSucceedsWithMock) {
+ MockHandleGuard
mg(reinterpret_cast<hdfsFile>(static_cast<uintptr_t>(0xdeadbeef)));
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ ExclusiveHdfsFileHandle handle(mock_fs, "/test/file.parquet", 12345);
+ ASSERT_TRUE(handle.init(4096).ok());
+ EXPECT_EQ(handle.file(), nullptr);
+
+ ASSERT_TRUE(handle.ensure_open().ok());
+ EXPECT_NE(handle.file(), nullptr);
+}
+
+// ensure_open() fails when mock returns nullptr.
+TEST(FileHandleCacheTest, EnsureOpenFailsWithMock) {
+ MockHandleGuard mg(nullptr);
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ ExclusiveHdfsFileHandle handle(mock_fs, "/test/file.parquet", 12345);
+ ASSERT_TRUE(handle.init(4096).ok());
+
+ auto st = handle.ensure_open();
+ ASSERT_FALSE(st.ok());
+ EXPECT_EQ(handle.file(), nullptr);
+}
+
+// ensure_open() is idempotent via call_once.
+TEST(FileHandleCacheTest, EnsureOpenIsIdempotentWithMock) {
+ MockHandleGuard
mg(reinterpret_cast<hdfsFile>(static_cast<uintptr_t>(0xdeadbeef)));
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ ExclusiveHdfsFileHandle handle(mock_fs, "/test/file.parquet", 12345);
+ ASSERT_TRUE(handle.init(4096).ok());
+
+ ASSERT_TRUE(handle.ensure_open().ok());
+ ASSERT_TRUE(handle.ensure_open().ok());
+}
+
+// init(-1) fetches file_size via mocked hdfsGetPathInfo.
+TEST(FileHandleCacheTest, InitWithUnknownFileSizeWithMock) {
+ MockHandleGuard mg(nullptr, 8192);
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ ExclusiveHdfsFileHandle handle(mock_fs, "/test/file.parquet", 12345);
+
+ ASSERT_TRUE(handle.init(-1).ok());
+ EXPECT_EQ(handle.file_size(), 8192);
+ EXPECT_EQ(handle.file(), nullptr);
+}
+
+// init(-1) fails when hdfsGetPathInfo returns nullptr.
+TEST(FileHandleCacheTest, InitFailsWhenGetPathInfoReturnsNull) {
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ ExclusiveHdfsFileHandle handle(mock_fs, "/test/file.parquet", 12345);
+
+ SyncPoint::get_instance()->enable_processing();
+ Defer defer {[&]() { SyncPoint::get_instance()->disable_processing(); }};
+ SyncPoint::CallbackGuard guard;
+ set_mock_return<hdfsFileInfo*>("HdfsFileHandle::init::hdfsGetPathInfo",
nullptr, &guard);
+ SyncPoint::CallbackGuard err_guard;
+ set_mock_return<std::string>("HdfsFileHandle::init::hdfs_error",
"connection refused",
+ &err_guard);
+
+ auto st = handle.init(-1);
+ ASSERT_FALSE(st.ok());
+ EXPECT_EQ(st.code(), TStatusCode::INTERNAL_ERROR);
+}
+
+// init(-1) returns NotFound when the file is missing.
+TEST(FileHandleCacheTest, InitReturnsNotFoundForMissingFile) {
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ ExclusiveHdfsFileHandle handle(mock_fs, "/test/missing.parquet", 12345);
+
+ SyncPoint::get_instance()->enable_processing();
+ Defer defer {[&]() { SyncPoint::get_instance()->disable_processing(); }};
+ SyncPoint::CallbackGuard guard;
+ set_mock_return<hdfsFileInfo*>("HdfsFileHandle::init::hdfsGetPathInfo",
nullptr, &guard);
+ SyncPoint::CallbackGuard err_guard;
+ set_mock_return<std::string>("HdfsFileHandle::init::hdfs_error", "No such
file or directory",
+ &err_guard);
+
+ auto st = handle.init(-1);
+ ASSERT_FALSE(st.ok());
+ EXPECT_EQ(st.code(), TStatusCode::NOT_FOUND);
+}
+
+// ensure_open() returns NotFound when error contains "No such file or
directory".
+TEST(FileHandleCacheTest, EnsureOpenReturnsNotFoundForMissingFile) {
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ ExclusiveHdfsFileHandle handle(mock_fs, "/test/missing.parquet", 12345);
+ ASSERT_TRUE(handle.init(4096).ok());
+
+ SyncPoint::get_instance()->enable_processing();
+ Defer defer {[&]() { SyncPoint::get_instance()->disable_processing(); }};
+ SyncPoint::CallbackGuard open_guard;
+ set_mock_return<hdfsFile>("HdfsFileHandle::ensure_open::hdfsOpenFile",
nullptr, &open_guard);
+ SyncPoint::CallbackGuard err_guard;
+ set_mock_return<std::string>("HdfsFileHandle::ensure_open::hdfs_error",
+ "No such file or directory", &err_guard);
+
+ auto st = handle.ensure_open();
+ ASSERT_FALSE(st.ok());
+ EXPECT_EQ(st.code(), TStatusCode::NOT_FOUND);
+}
+
+// --- Cache lifecycle tests ---
+
+// Helper: create a FileHandleCache with small capacity for testing.
+static std::unique_ptr<FileHandleCache> make_test_cache() {
+ return std::make_unique<FileHandleCache>(4, 1, 0);
+}
+
+// Helper: get a file handle from cache, asserting success.
+static void get_handle(FileHandleCache& cache, const hdfsFS& fs, const
std::string& fname,
+ int64_t mtime, FileHandleCache::Accessor* accessor,
bool* cache_hit) {
+ ASSERT_TRUE(cache.get_file_handle(fs, fname, mtime, 4096, false, accessor,
cache_hit).ok());
+}
+
+// ensure_open succeeds → ~Accessor releases (unbuffer OK) → next
get_file_handle hits cache.
+TEST(FileHandleCacheTest, OpenedHandleReleasedBackToCache) {
+ MockHandleGuard
mg(reinterpret_cast<hdfsFile>(static_cast<uintptr_t>(0xdeadbeef)));
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ auto cache = make_test_cache();
+ const std::string fname = "/test/opened_release.parquet";
+ constexpr int64_t mtime = 12345;
+
+ bool cache_hit = false;
+ {
+ FileHandleCache::Accessor accessor;
+ get_handle(*cache, mock_fs, fname, mtime, &accessor, &cache_hit);
+ EXPECT_FALSE(cache_hit);
+ ASSERT_TRUE(accessor.get()->ensure_open().ok());
+ EXPECT_NE(accessor.get()->file(), nullptr);
+ }
+
+ FileHandleCache::Accessor accessor2;
+ get_handle(*cache, mock_fs, fname, mtime, &accessor2, &cache_hit);
+#ifdef USE_HADOOP_HDFS
+ EXPECT_TRUE(cache_hit);
+ EXPECT_NE(accessor2.get()->file(), nullptr);
+#else
+ // libhdfs3: ~Accessor() destroys the opened handle, so the next lookup
must miss.
+ EXPECT_FALSE(cache_hit);
+ EXPECT_EQ(accessor2.get()->file(), nullptr);
+#endif
+}
+
+// ensure_open fails → ~Accessor destroys → next get_file_handle misses cache.
+TEST(FileHandleCacheTest, OpenFailedHandleDestroyedNotCached) {
+ MockHandleGuard mg(nullptr);
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ auto cache = make_test_cache();
+ const std::string fname = "/test/open_fail.parquet";
+ constexpr int64_t mtime = 12345;
+
+ bool cache_hit = false;
+ {
+ FileHandleCache::Accessor accessor;
+ get_handle(*cache, mock_fs, fname, mtime, &accessor, &cache_hit);
+ EXPECT_FALSE(cache_hit);
+ auto st = accessor.get()->ensure_open();
+ ASSERT_FALSE(st.ok());
+ EXPECT_EQ(accessor.get()->file(), nullptr);
+ }
+
+ FileHandleCache::Accessor accessor2;
+ get_handle(*cache, mock_fs, fname, mtime, &accessor2, &cache_hit);
+ EXPECT_FALSE(cache_hit);
+}
+
+// read_at triggers ensure_open failure → reader destroyed → ~Accessor
destroys → cache miss.
+TEST(FileHandleCacheTest, ReadAtOpenFailedHandleDestroyedNotCached) {
+ MockHandleGuard mg(nullptr);
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ auto cache = make_test_cache();
+ const std::string fname = "/test/read_at_fail.parquet";
+ constexpr int64_t mtime = 12345;
+
+ bool cache_hit = false;
+ {
+ FileHandleCache::Accessor accessor;
+ get_handle(*cache, mock_fs, fname, mtime, &accessor, &cache_hit);
+ EXPECT_FALSE(cache_hit);
+ auto reader =
+ std::make_shared<HdfsFileReader>(Path(fname), "hdfs",
std::move(accessor), mtime);
+ char buf[16];
+ size_t bytes_read = 0;
+ auto st = reader->read_at(0, {buf, sizeof(buf)}, &bytes_read, nullptr);
+ ASSERT_FALSE(st.ok());
+ }
+
+ FileHandleCache::Accessor accessor2;
+ get_handle(*cache, mock_fs, fname, mtime, &accessor2, &cache_hit);
+ EXPECT_FALSE(cache_hit);
+}
+
+// Serial ensure_open() preserves Status inside call_once.
+TEST(FileHandleCacheTest, OpenFailurePreservesStatusInsideCallOnce) {
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ ExclusiveHdfsFileHandle handle(mock_fs, "/test/serial_fail.parquet",
12345);
+ ASSERT_TRUE(handle.init(4096).ok());
+
+ auto* sp = SyncPoint::get_instance();
+ sp->enable_processing();
+ Defer defer {[&]() { sp->disable_processing(); }};
+ SyncPoint::CallbackGuard open_guard;
+ set_mock_return<hdfsFile>("HdfsFileHandle::ensure_open::hdfsOpenFile",
nullptr, &open_guard);
+ SyncPoint::CallbackGuard err_guard;
+ set_mock_return<std::string>("HdfsFileHandle::ensure_open::hdfs_error",
+ "No such file or directory", &err_guard);
+
+ // First call opens and fails -> NotFound.
+ auto st1 = handle.ensure_open();
+ ASSERT_FALSE(st1.ok());
+ EXPECT_EQ(st1.code(), TStatusCode::NOT_FOUND);
+
+ // Change hdfs_error mock to a different value; _open_status must be
preserved.
+ sp->clear_call_back("HdfsFileHandle::ensure_open::hdfs_error");
+ set_mock_return<std::string>("HdfsFileHandle::ensure_open::hdfs_error",
"Permission denied",
+ &err_guard);
+
+ auto st2 = handle.ensure_open();
+ ASSERT_FALSE(st2.ok());
+ EXPECT_EQ(st2.code(), TStatusCode::NOT_FOUND);
+}
+
+// Concurrent ensure_open() must return the same Status to all callers.
+TEST(FileHandleCacheTest, ConcurrentOpenFailureReturnsSameStatusToAllCallers) {
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ ExclusiveHdfsFileHandle handle(mock_fs, "/test/concurrent_fail.parquet",
12345);
+ ASSERT_TRUE(handle.init(4096).ok());
+
+ auto* sp = SyncPoint::get_instance();
+ sp->enable_processing();
+ Defer defer {[&]() { sp->disable_processing(); }};
+ SyncPoint::CallbackGuard open_guard;
+ set_mock_return<hdfsFile>("HdfsFileHandle::ensure_open::hdfsOpenFile",
nullptr, &open_guard);
+
+ // Simulate libhdfs thread-local last-error: only first call returns real
message.
+ std::atomic<int> err_call_count {0};
+ SyncPoint::CallbackGuard err_guard;
+ sp->set_call_back(
+ "HdfsFileHandle::ensure_open::hdfs_error",
+ [&err_call_count](auto&& args) {
+ auto* ret = try_any_cast_ret<std::string>(args);
+ if (err_call_count.fetch_add(1) == 0) {
+ ret->first = "No such file or directory";
+ } else {
+ ret->first = "";
+ }
+ ret->second = true;
+ },
+ &err_guard);
+
+ constexpr int kNumThreads = 8;
+ std::vector<std::thread> threads;
+ std::vector<TStatusCode::type> codes(kNumThreads);
+ for (int i = 0; i < kNumThreads; ++i) {
+ threads.emplace_back([&handle, &codes, i]() {
+ auto st = handle.ensure_open();
+ codes[i] = static_cast<TStatusCode::type>(st.code());
+ });
+ }
+ for (auto& t : threads) {
+ t.join();
+ }
+
+ for (int i = 0; i < kNumThreads; ++i) {
+ EXPECT_EQ(codes[i], TStatusCode::NOT_FOUND) << "thread " << i << " got
code " << codes[i];
+ }
+}
+
+// A read after close must fail before lazy open: no hdfsOpenFile call is
allowed.
+TEST(FileHandleCacheTest, ReadAfterCloseSkipsLazyOpen) {
+ MockHandleGuard
mg(reinterpret_cast<hdfsFile>(static_cast<uintptr_t>(0xdeadbeef)));
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ auto cache = make_test_cache();
+ const std::string fname = "/test/close_before_read.parquet";
+ constexpr int64_t mtime = 12345;
+
+ bool cache_hit = false;
+ FileHandleCache::Accessor accessor;
+ get_handle(*cache, mock_fs, fname, mtime, &accessor, &cache_hit);
+ auto reader = std::make_shared<HdfsFileReader>(Path(fname), "hdfs",
std::move(accessor), mtime);
+
+ ASSERT_TRUE(reader->close().ok());
+ char buf[16];
+ size_t bytes_read = 0;
+ auto st = reader->read_at(0, {buf, sizeof(buf)}, &bytes_read, nullptr);
+ ASSERT_FALSE(st.ok());
+ EXPECT_EQ(st.code(), TStatusCode::INTERNAL_ERROR);
+ // "read closed file" proves the closed guard fired before lazy open.
+ EXPECT_NE(st.to_string().find("read closed file"), std::string::npos);
+}
+
+// Second read_at after a failed read must not dereference null _handle.
+TEST(FileHandleCacheTest, SecondReadAfterFailureDoesNotCrash) {
+ MockHandleGuard
mg(reinterpret_cast<hdfsFile>(static_cast<uintptr_t>(0xdeadbeef)));
+ auto mock_fs = reinterpret_cast<hdfsFS>(static_cast<uintptr_t>(0x1));
+ auto cache = make_test_cache();
+ const std::string fname = "/test/second_read.parquet";
+ constexpr int64_t mtime = 12345;
+
+ bool cache_hit = false;
+ FileHandleCache::Accessor accessor;
+ get_handle(*cache, mock_fs, fname, mtime, &accessor, &cache_hit);
+ auto reader = std::make_shared<HdfsFileReader>(Path(fname), "hdfs",
std::move(accessor), mtime);
+
+ char buf[16];
+ size_t bytes_read = 0;
+ // offset > file_size(4096) -> IOError -> read_at_impl sets
_handle=nullptr.
+ auto st1 = reader->read_at(5000, {buf, sizeof(buf)}, &bytes_read, nullptr);
+ ASSERT_FALSE(st1.ok());
+ // Confirm error came from do_read_at_impl's offset guard (sets
_handle=nullptr).
+ EXPECT_EQ(st1.code(), TStatusCode::IO_ERROR);
+
+ // Second read must not crash; should return InternalError about destroyed
handle.
+ auto st2 = reader->read_at(0, {buf, sizeof(buf)}, &bytes_read, nullptr);
+ ASSERT_FALSE(st2.ok());
+ EXPECT_EQ(st2.code(), TStatusCode::INTERNAL_ERROR);
+}
+
} // namespace doris::io
diff --git a/be/test/io/fs/hdfs_file_system_test.cpp
b/be/test/io/fs/hdfs_file_system_test.cpp
index db443951a26..5d6cef97f2e 100644
--- a/be/test/io/fs/hdfs_file_system_test.cpp
+++ b/be/test/io/fs/hdfs_file_system_test.cpp
@@ -15,14 +15,22 @@
// specific language governing permissions and limitations
// under the License.
+#include "io/fs/hdfs_file_system.h"
+
#include <gtest/gtest.h>
+#include <map>
+#include <string>
+
#include "common/config.h"
#include "cpp/sync_point.h"
+#include "gen_cpp/PlanNodes_types.h"
+#include "gen_cpp/Status_types.h"
#include "io/fs/file_reader.h"
#include "io/fs/file_writer.h"
#include "io/fs/hdfs_file_writer.h"
#include "io/fs/local_file_system.h"
+#include "util/defer_op.h"
namespace doris {
@@ -151,4 +159,66 @@ TEST(HdfsFileSystemTest, Write) {
st = local_fs->delete_directory(test_dir);
}
+// Guarded: the enable_java_support check is compiled out on libhdfs3 builds.
+#ifdef USE_HADOOP_HDFS
+// create() returns error when java support is disabled.
+TEST(HdfsFileSystemTest, CreateFailsWhenJavaSupportDisabled) {
+ const bool old_enable_java_support = config::enable_java_support;
+ config::enable_java_support = false;
+ Defer defer {[&]() { config::enable_java_support =
old_enable_java_support; }};
+ std::map<std::string, std::string> properties;
+ auto res = io::HdfsFileSystem::create(properties, "hdfs://namenode:8020",
"test_id", "/");
+ ASSERT_FALSE(res.has_value());
+ EXPECT_NE(res.error().to_string().find("enable_java_support"),
std::string::npos);
+}
+#endif
+
+// open_file_internal returns IOError when _fs_handler is null.
+TEST(HdfsFileSystemTest, OpenFileFailsWithoutHandler) {
+ THdfsParams params;
+ io::HdfsFileSystem fs(params, "hdfs://namenode:8020", "test_id", "/");
+ io::FileReaderSPtr reader;
+ auto st = fs.open_file_internal("/test/file.parquet", &reader,
io::FileReaderOptions::DEFAULT);
+ ASSERT_FALSE(st.ok());
+ EXPECT_EQ(st.code(), TStatusCode::IO_ERROR);
+}
+
+// file_size_impl returns IOError when _fs_handler is null.
+TEST(HdfsFileSystemTest, FileSizeFailsWithoutHandler) {
+ THdfsParams params;
+ io::HdfsFileSystem fs(params, "hdfs://namenode:8020", "test_id", "/");
+ int64_t file_size = 0;
+ auto st = fs.file_size_impl("/test/file.parquet", &file_size);
+ ASSERT_FALSE(st.ok());
+ EXPECT_EQ(st.code(), TStatusCode::IO_ERROR);
+}
+
+// exists_impl returns IOError when _fs_handler is null.
+TEST(HdfsFileSystemTest, ExistsFailsWithoutHandler) {
+ THdfsParams params;
+ io::HdfsFileSystem fs(params, "hdfs://namenode:8020", "test_id", "/");
+ bool exists = false;
+ auto st = fs.exists_impl("/test/file.parquet", &exists);
+ ASSERT_FALSE(st.ok());
+ EXPECT_EQ(st.code(), TStatusCode::IO_ERROR);
+}
+
+// create_directory_impl returns IOError when _fs_handler is null.
+TEST(HdfsFileSystemTest, CreateDirectoryFailsWithoutHandler) {
+ THdfsParams params;
+ io::HdfsFileSystem fs(params, "hdfs://namenode:8020", "test_id", "/");
+ auto st = fs.create_directory_impl("/test/dir", false);
+ ASSERT_FALSE(st.ok());
+ EXPECT_EQ(st.code(), TStatusCode::IO_ERROR);
+}
+
+// rename_impl returns IOError when _fs_handler is null.
+TEST(HdfsFileSystemTest, RenameFailsWithoutHandler) {
+ THdfsParams params;
+ io::HdfsFileSystem fs(params, "hdfs://namenode:8020", "test_id", "/");
+ auto st = fs.rename_impl("/test/old.parquet", "/test/new.parquet");
+ ASSERT_FALSE(st.ok());
+ EXPECT_EQ(st.code(), TStatusCode::IO_ERROR);
+}
+
} // namespace doris
diff --git a/be/test/storage/index/snii/snii_blob_directory_test.cpp
b/be/test/storage/index/snii/snii_blob_directory_test.cpp
index b5988750c87..f6b194a2872 100644
--- a/be/test/storage/index/snii/snii_blob_directory_test.cpp
+++ b/be/test/storage/index/snii/snii_blob_directory_test.cpp
@@ -39,10 +39,13 @@
#include <string>
#include <vector>
+#include "common/config.h"
#include "common/status.h"
#include "io/fs/local_file_system.h"
#include "storage/index/snii/format/metadata_directory.h"
#include "storage/index/snii/snii_doris_adapter.h"
+#include "util/debug_points.h"
+#include "util/defer_op.h"
using doris::Status;
using doris::segment_v2::snii_doris::DorisSniiFileReader;
@@ -209,6 +212,39 @@ TEST(SniiBlobDirectory, MissingFileFailsWithoutThrowing) {
EXPECT_EQ(CL_ERR_IO, err.number());
}
+// A read-stage NotFound must surface as CL_ERR_FileNotFound, not collapse
into CL_ERR_IO.
+TEST(SniiBlobDirectory, ReadFailureKeepsNotFoundDistinguishable) {
+ Fixture fx;
+ fx.Build();
+ ASSERT_FALSE(::testing::Test::HasFatalFailure());
+
+ const auto old_enable = doris::config::enable_debug_points;
+ doris::config::enable_debug_points = true;
+ const std::string point = "LocalFileReader::read_at_impl.io_error";
+ doris::DebugPoints::instance()->add_with_params(
+ point, {{"errno", std::to_string(ENOENT)}, {"sub_path",
"snii_blob_dir_test"}});
+ doris::Defer dp {[&]() {
+ doris::DebugPoints::instance()->remove(point);
+ doris::config::enable_debug_points = old_enable;
+ }};
+
+ lucene::store::IndexInput* raw = nullptr;
+ CLuceneError err;
+ ASSERT_TRUE(fx.dir->openInput("a", raw, err));
+ std::unique_ptr<lucene::store::IndexInput> input(raw);
+
+ std::vector<uint8_t> got(fx.payload_a.size(), 0);
+ bool threw = false;
+ try {
+ input->readBytes(got.data(), static_cast<int32_t>(got.size()));
+ } catch (CLuceneError& e) {
+ threw = true;
+ EXPECT_EQ(CL_ERR_FileNotFound, e.number());
+ }
+ EXPECT_TRUE(threw);
+ input->close();
+}
+
TEST(SniiBlobDirectory, WriteOperationsThrowUnsupportedAndCloseNeverThrows) {
Fixture fx;
fx.Build();
diff --git a/be/test/storage/segment/inverted_index_file_reader_test.cpp
b/be/test/storage/segment/inverted_index_file_reader_test.cpp
index 5a5a2633a15..e4605304565 100644
--- a/be/test/storage/segment/inverted_index_file_reader_test.cpp
+++ b/be/test/storage/segment/inverted_index_file_reader_test.cpp
@@ -16,6 +16,7 @@
// under the License.
#include <CLucene.h>
+#include <errno.h>
#include <gtest/gtest.h>
#include <unistd.h>
@@ -24,6 +25,7 @@
#include <string>
#include <vector>
+#include "common/config.h"
#include "io/fs/local_file_system.h"
#include "runtime/exec_env.h"
#include "storage/data_dir.h"
@@ -39,6 +41,7 @@
#include "storage/options.h"
#include "storage/storage_engine.h"
#include "storage/tablet/tablet_schema.h"
+#include "util/debug_points.h"
namespace doris::segment_v2 {
@@ -297,6 +300,30 @@ TEST_F(InvertedIndexFileReaderTest,
TestUnknownIndexFormatError) {
status.msg().find("CLuceneError") != std::string::npos);
}
+// A read-stage NotFound must surface as INVERTED_INDEX_FILE_NOT_FOUND so
callers can downgrade.
+TEST_F(InvertedIndexFileReaderTest, TestV2ReadNotFoundReturnsFileNotFound) {
+ std::string index_path = kTestDir + "/read_not_found_index_file";
+ create_invalid_version_file(index_path + ".idx", 1);
+
+ InvertedIndexFileInfo file_info;
+ file_info.set_index_size(1024); // size known: init() skips stat, open
succeeds
+
+ IndexFileReader reader(io::global_local_filesystem(), index_path,
+ InvertedIndexStorageFormatPB::V2, file_info);
+
+ const auto old_enable = config::enable_debug_points;
+ config::enable_debug_points = true;
+ const std::string point = "LocalFileReader::read_at_impl.io_error";
+ DebugPoints::instance()->add_with_params(
+ point, {{"errno", std::to_string(ENOENT)}, {"sub_path",
"read_not_found"}});
+ Status status = reader.init(4096);
+ DebugPoints::instance()->remove(point);
+ config::enable_debug_points = old_enable;
+
+ EXPECT_FALSE(status.ok());
+ EXPECT_EQ(status.code(), ErrorCode::INVERTED_INDEX_FILE_NOT_FOUND);
+}
+
// Test case for V1 format file not found error
TEST_F(InvertedIndexFileReaderTest, TestV1FileNotFoundError) {
std::string index_path = kTestDir + "/non_existent";
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 eadfaaa484b..b04b5ec82e0 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
@@ -1746,13 +1746,13 @@ public class IcebergScanPlanProvider implements
ConnectorScanPlanProvider {
validateDeletionVectorMetadata(
delete.path().toString(), delete.fileSizeInBytes(),
contentOffset, contentLength);
return IcebergScanRange.DeleteFile.deletionVector(path,
lowerBound, upperBound,
- contentOffset, contentLength);
+ contentOffset, contentLength,
delete.fileSizeInBytes());
}
return IcebergScanRange.DeleteFile.positionDelete(path,
deleteFileFormat(delete.format()),
- lowerBound, upperBound);
+ lowerBound, upperBound, delete.fileSizeInBytes());
} else if (content == FileContent.EQUALITY_DELETES) {
return IcebergScanRange.DeleteFile.equalityDelete(path,
deleteFileFormat(delete.format()),
- delete.equalityFieldIds());
+ delete.equalityFieldIds(), delete.fileSizeInBytes());
}
// Defensive (legacy parity): delete files are only position or
equality; DATA content here is a bug.
throw new IllegalStateException("Unknown delete content: " + content);
diff --git
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanRange.java
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanRange.java
index 0e9f3548831..6a4a6316830 100644
---
a/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanRange.java
+++
b/fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergScanRange.java
@@ -348,6 +348,9 @@ public class IcebergScanRange implements ConnectorScanRange
{
deleteFileDesc.setOriginalPath(positionDeleteOriginalPath);
deleteFileDesc.setFileFormat(positionDeleteFileFormat);
deleteFileDesc.setContent(positionDeleteContent);
+ if (rangeDesc.isSetFileSize()) {
+ deleteFileDesc.setFileSize(rangeDesc.getFileSize());
+ }
if (positionDeleteContentOffset != null) {
deleteFileDesc.setContentOffset(positionDeleteContentOffset);
}
@@ -631,9 +634,11 @@ public class IcebergScanRange implements
ConnectorScanRange {
// deletion vector only (null otherwise).
private final Long contentOffset;
private final Long contentSizeInBytes;
+ private final long fileSize;
private DeleteFile(String path, int content, TFileFormatType
fileFormat, Long positionLowerBound,
- Long positionUpperBound, List<Integer> fieldIds, Long
contentOffset, Long contentSizeInBytes) {
+ Long positionUpperBound, List<Integer> fieldIds, Long
contentOffset, Long contentSizeInBytes,
+ long fileSize) {
this.path = path;
this.content = content;
this.fileFormat = fileFormat;
@@ -642,13 +647,14 @@ public class IcebergScanRange implements
ConnectorScanRange {
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 +663,16 @@ public class IcebergScanRange implements
ConnectorScanRange {
* 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<Integer> fieldIds) {
- return new DeleteFile(path, CONTENT_EQUALITY_DELETE, fileFormat,
null, null, fieldIds, null, null);
+ public static DeleteFile equalityDelete(String path, TFileFormatType
fileFormat, List<Integer> fieldIds,
+ long fileSize) {
+ return new DeleteFile(path, CONTENT_EQUALITY_DELETE, fileFormat,
null, null,
+ fieldIds, null, null, fileSize);
}
int getContent() {
@@ -693,6 +701,7 @@ public class IcebergScanRange implements ConnectorScanRange
{
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 b8509baf2ee..512c2ffa338 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
@@ -4292,4 +4292,112 @@ public class IcebergScanPlanProviderTest {
throw new UnsupportedOperationException();
}
}
+
+ @Test
+ public void deleteFileSizePropagatedFromIcebergManifest() {
+ 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();
+ 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<ConnectorScanRange> 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);
+
+ List<TIcebergDeleteFileDesc> descs = range.rewritableDeleteDescs();
+ Assertions.assertFalse(descs.isEmpty());
+ TIcebergDeleteFileDesc desc = descs.get(0);
+ Assertions.assertTrue(desc.isSetFileSize());
+ Assertions.assertEquals(128L, desc.getFileSize());
+ }
+
+ @Test
+ public void deletionVectorFileSizePropagatedFromIcebergManifest() {
+ Table table = createTable("t_dv_filesize", SCHEMA,
PartitionSpec.unpartitioned(),
+ Collections.singletonMap("format-version", "3"));
+ table.newAppend()
+ .appendFile(dataFile(table.spec(),
"s3://b/db/t_dv_filesize/f1.parquet", 512, null, null))
+ .commit();
+ DeleteFile dvDelete = FileMetadata.deleteFileBuilder(table.spec())
+ .ofPositionDeletes()
+ .withPath("s3://b/db/t_dv_filesize/dv.puffin")
+ .withFormat(FileFormat.PUFFIN)
+ .withFileSizeInBytes(256L)
+ .withRecordCount(1L)
+ .withContentOffset(16L)
+ .withContentSizeInBytes(64L)
+ .withReferencedDataFile("s3://b/db/t_dv_filesize/f1.parquet")
+ .build();
+ table.newRowDelta().addDeletes(dvDelete).commit();
+
+ IcebergScanPlanProvider provider = new IcebergScanPlanProvider(
+ IcebergCatalogProperties.of(Collections.emptyMap()),
opsReturning(table));
+ List<ConnectorScanRange> ranges = provider.planScan(
+ new FakeScanSession("UTC", Collections.emptyMap()),
+ ConnectorScanRequest.builder(new IcebergTableHandle("db1",
"t_dv_filesize"), Collections.emptyList())
+ .build());
+
+ Assertions.assertEquals(1, ranges.size());
+ IcebergScanRange range = (IcebergScanRange) ranges.get(0);
+
+ List<TIcebergDeleteFileDesc> descs = range.rewritableDeleteDescs();
+ Assertions.assertFalse(descs.isEmpty());
+ TIcebergDeleteFileDesc desc = descs.get(0);
+ Assertions.assertEquals(3, desc.getContent());
+ Assertions.assertTrue(desc.isSetFileSize());
+ Assertions.assertEquals(256L, desc.getFileSize());
+ }
+
+ @Test
+ public void equalityDeleteFileSizePropagatedFromIcebergManifest() {
+ Table table = createTable("t_eq_filesize", SCHEMA,
PartitionSpec.unpartitioned(),
+ Collections.singletonMap("format-version", "2"));
+ table.newAppend()
+ .appendFile(dataFile(table.spec(),
"s3://b/db/t_eq_filesize/f1.parquet", 512, null, null))
+ .commit();
+ DeleteFile eqDelete = FileMetadata.deleteFileBuilder(table.spec())
+ .ofEqualityDeletes(1)
+ .withPath("s3://b/db/t_eq_filesize/eq.parquet")
+ .withFormat(FileFormat.PARQUET)
+ .withFileSizeInBytes(96L)
+ .withRecordCount(1L)
+ .build();
+ table.newRowDelta().addDeletes(eqDelete).commit();
+
+ IcebergScanPlanProvider provider = new IcebergScanPlanProvider(
+ IcebergCatalogProperties.of(Collections.emptyMap()),
opsReturning(table));
+ List<ConnectorScanRange> ranges = provider.planScan(
+ new FakeScanSession("UTC", Collections.emptyMap()),
+ ConnectorScanRequest.builder(new IcebergTableHandle("db1",
"t_eq_filesize"), Collections.emptyList())
+ .build());
+
+ Assertions.assertEquals(1, ranges.size());
+ IcebergScanRange range = (IcebergScanRange) ranges.get(0);
+
+ // equality deletes are excluded from rewritableDeleteDescs; verify
via populateRangeParams
+ TFileRangeDesc rangeDesc = new TFileRangeDesc();
+ TTableFormatFileDesc formatDesc = new TTableFormatFileDesc();
+ range.populateRangeParams(formatDesc, rangeDesc);
+ List<TIcebergDeleteFileDesc> descs =
formatDesc.getIcebergParams().getDeleteFiles();
+ Assertions.assertFalse(descs.isEmpty());
+ TIcebergDeleteFileDesc desc = descs.get(0);
+ Assertions.assertEquals(2, desc.getContent());
+ Assertions.assertTrue(desc.isSetFileSize());
+ Assertions.assertEquals(96L, 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 07f52b5d3a5..2902283d70d 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 class IcebergScanRangeTest {
// 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();
@@ -207,6 +207,9 @@ public class IcebergScanRangeTest {
Assertions.assertEquals(10L, d.getPositionLowerBound());
Assertions.assertTrue(d.isSetPositionUpperBound());
Assertions.assertEquals(99L, d.getPositionUpperBound());
+ // file_size is propagated from DeleteFile to TIcebergDeleteFileDesc.
+ Assertions.assertTrue(d.isSetFileSize());
+ Assertions.assertEquals(100L, d.getFileSize());
// A position delete carries neither equality field-ids nor a
deletion-vector blob ref.
Assertions.assertFalse(d.isSetFieldIds());
Assertions.assertFalse(d.isSetContentOffset());
@@ -218,7 +221,7 @@ public class IcebergScanRangeTest {
// 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 +239,9 @@ public class IcebergScanRangeTest {
// 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();
@@ -255,6 +258,8 @@ public class IcebergScanRangeTest {
Assertions.assertEquals(64L, dvDesc.getContentSizeInBytes());
Assertions.assertEquals(5L, dvDesc.getPositionLowerBound());
Assertions.assertEquals(42L, dvDesc.getPositionUpperBound());
+ Assertions.assertTrue(dvDesc.isSetFileSize());
+ Assertions.assertEquals(100L, dvDesc.getFileSize());
TIcebergDeleteFileDesc eqDesc = deletes.get(1);
Assertions.assertEquals(2, eqDesc.getContent());
@@ -262,6 +267,8 @@ public class IcebergScanRangeTest {
Assertions.assertEquals(TFileFormatType.FORMAT_PARQUET,
eqDesc.getFileFormat());
Assertions.assertFalse(eqDesc.isSetContentOffset());
Assertions.assertFalse(eqDesc.isSetPositionLowerBound());
+ Assertions.assertTrue(eqDesc.isSetFileSize());
+ Assertions.assertEquals(100L, eqDesc.getFileSize());
}
@Test
@@ -446,6 +453,31 @@ public class IcebergScanRangeTest {
Assertions.assertFalse(rangeDesc.isSetColumnsFromPath());
}
+ @Test
+ public void populateRangeParamsSysTableForwardsDeleteFileSize() {
+ // The sys-table delete descriptor must forward the range's file_size;
BE keys on delete_files[0].file_size.
+ IcebergScanRange range = new IcebergScanRange.Builder()
+ .path("/puffin/delete-file.puffin")
+ .positionDeleteSysTableSplit(3,
TFileFormatType.FORMAT_PARQUET, "/raw/delete.puffin")
+ .positionDeleteDeletionVector("/data/file.parquet", 4L, 100L)
+ .build();
+
+ TFileRangeDesc rangeDesc = new TFileRangeDesc();
+ rangeDesc.setFileSize(12345L);
+ populate(range, rangeDesc);
+
+ TIcebergFileDesc fd =
rangeDesc.getTableFormatParams().getIcebergParams();
+ Assertions.assertTrue(fd.isSetDeleteFiles());
+ Assertions.assertEquals(1, fd.getDeleteFilesSize());
+ Assertions.assertEquals(12345L,
fd.getDeleteFiles().get(0).getFileSize());
+
+ // isSet guard: an unset range size must not emit file_size on the
delete descriptor.
+ TFileRangeDesc unsetDesc = populate(range, new TFileRangeDesc());
+ Assertions.assertFalse(
+
unsetDesc.getTableFormatParams().getIcebergParams().getDeleteFiles().get(0)
+ .isSetFileSize());
+ }
+
// ── commit-bridge supply (S4 part 2): getOriginalPath() key +
rewritableDeleteDescs() non-equality filter ──
@Test
@@ -467,11 +499,11 @@ public class IcebergScanRangeTest {
// 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 +527,34 @@ public class IcebergScanRangeTest {
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 9867ab16052..1c30a652fe9 100644
--- a/gensrc/thrift/PlanNodes.thrift
+++ b/gensrc/thrift/PlanNodes.thrift
@@ -333,6 +333,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 {
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]