Skip to content

Commit d4e3c49

Browse files
committed
[fix](io) Fix iceberg delete file file_size propagation and add HDFS read path bvar
1 parent d6ca1a6 commit d4e3c49

7 files changed

Lines changed: 40 additions & 30 deletions

File tree

be/src/format/table/iceberg_delete_file_reader_helper.cpp

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -241,8 +241,7 @@ TFileRangeDesc build_iceberg_delete_file_range(const std::string& path, int64_t
241241
range.path = path;
242242
range.start_offset = 0;
243243
range.size = -1;
244-
// thrift optional defaults to 0 when unset; treat 0 as unknown (-1)
245-
range.file_size = (file_size <= 0) ? -1 : file_size;
244+
range.file_size = file_size;
246245
return range;
247246
}
248247

@@ -277,8 +276,9 @@ Status read_iceberg_position_delete_file(const TIcebergDeleteFileDesc& delete_fi
277276
return Status::InvalidArgument("invalid position delete reader options");
278277
}
279278

280-
TFileRangeDesc delete_range =
281-
build_iceberg_delete_file_range(delete_file.path, delete_file.file_size);
279+
TFileRangeDesc delete_range = build_iceberg_delete_file_range(
280+
delete_file.path,
281+
delete_file.__isset.file_size ? delete_file.file_size : -1);
282282
if (options.fs_name != nullptr && !options.fs_name->empty()) {
283283
delete_range.__set_fs_name(*options.fs_name);
284284
}
@@ -350,8 +350,9 @@ Status read_iceberg_deletion_vector(const TIcebergDeleteFileDesc& delete_file,
350350
DBUG_EXECUTE_IF("IcebergDeleteFileReader.read_deletion_vector.should_stop",
351351
{ return Status::EndOfFile("stop read."); });
352352

353-
TFileRangeDesc delete_range =
354-
build_iceberg_delete_file_range(delete_file.path, delete_file.file_size);
353+
TFileRangeDesc delete_range = build_iceberg_delete_file_range(
354+
delete_file.path,
355+
delete_file.__isset.file_size ? delete_file.file_size : -1);
355356
if (options.fs_name != nullptr && !options.fs_name->empty()) {
356357
delete_range.__set_fs_name(*options.fs_name);
357358
}

be/src/format_v2/table/iceberg_reader.cpp

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1165,7 +1165,7 @@ Status IcebergTableReader::_parse_deletion_vector_file(const TTableFormatFileDes
11651165
desc->path = deletion_vector->path;
11661166
desc->start_offset = deletion_vector->content_offset;
11671167
desc->size = static_cast<int64_t>(bytes_read);
1168-
desc->file_size = deletion_vector->file_size;
1168+
desc->file_size = deletion_vector->__isset.file_size ? deletion_vector->file_size : -1;
11691169
desc->format = DeleteFileDesc::Format::ICEBERG;
11701170
*has_delete_file = true;
11711171
return Status::OK();
@@ -1521,7 +1521,9 @@ Status IcebergTableReader::_create_delete_file_reader(const TIcebergDeleteFileDe
15211521
return Status::NotSupported("Unsupported Iceberg delete file format {}",
15221522
delete_file.file_format);
15231523
}
1524-
auto delete_range = build_iceberg_delete_file_range(delete_file.path, delete_file.file_size);
1524+
auto delete_range = build_iceberg_delete_file_range(
1525+
delete_file.path,
1526+
delete_file.__isset.file_size ? delete_file.file_size : -1);
15251527
if (_current_task != nullptr && _current_task->data_file != nullptr &&
15261528
!_current_task->data_file->fs_name.empty()) {
15271529
delete_range.__set_fs_name(_current_task->data_file->fs_name);

be/src/io/fs/file_handle_cache.cpp

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,13 +27,16 @@
2727

2828
#include "common/cast_set.h"
2929
#include "io/fs/err_utils.h"
30+
#include "io/hdfs_util.h"
31+
#include "util/bvar_helper.h"
3032
#include "util/hash_util.hpp"
3133
#include "util/time.h"
3234
namespace doris::io {
3335

3436
HdfsFileHandle::~HdfsFileHandle() {
3537
if (_hdfs_file != nullptr && _fs != nullptr) {
3638
VLOG_FILE << "hdfsCloseFile() fid=" << _hdfs_file;
39+
SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_close_latency);
3740
hdfsCloseFile(_fs, _hdfs_file); // TODO: check return code
3841
}
3942
_fs = nullptr;
@@ -57,6 +60,7 @@ Status HdfsFileHandle::init(int64_t file_size) {
5760
Status HdfsFileHandle::ensure_open() {
5861
std::call_once(_open_once, [this]() {
5962
VLOG_DEBUG << "lazy open hdfs file: " << _fname;
63+
SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_open_latency);
6064
_hdfs_file = hdfsOpenFile(_fs, _fname.c_str(), O_RDONLY, 0, 0, 0);
6165
});
6266
if (_hdfs_file == nullptr) {

be/src/io/fs/hdfs_file_reader.cpp

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@
3838
#include "runtime/workload_management/io_throttle.h"
3939
#include "runtime/workload_management/resource_context.h"
4040
#include "service/backend_options.h"
41+
#include "util/bvar_helper.h"
4142

4243
namespace doris::io {
4344

@@ -164,8 +165,12 @@ Status HdfsFileReader::do_read_at_impl(size_t offset, Slice result, size_t* byte
164165
int64_t max_to_read = bytes_req - has_read;
165166
tSize to_read = static_cast<tSize>(
166167
std::min(max_to_read, static_cast<int64_t>(std::numeric_limits<tSize>::max())));
167-
tSize loop_read = hdfsPread(_handle->fs(), _handle->file(), offset + has_read,
168-
to + has_read, to_read);
168+
tSize loop_read;
169+
{
170+
SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_read_latency);
171+
loop_read = hdfsPread(_handle->fs(), _handle->file(), offset + has_read,
172+
to + has_read, to_read);
173+
}
169174
{
170175
[[maybe_unused]] Status error_ret;
171176
TEST_INJECTION_POINT_RETURN_WITH_VALUE("HdfsFileReader:read_error", error_ret);
@@ -230,8 +235,12 @@ Status HdfsFileReader::do_read_at_impl(size_t offset, Slice result, size_t* byte
230235

231236
size_t has_read = 0;
232237
while (has_read < bytes_req) {
233-
int64_t loop_read = hdfsRead(_handle->fs(), _handle->file(), to + has_read,
234-
static_cast<int32_t>(bytes_req - has_read));
238+
int64_t loop_read;
239+
{
240+
SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_read_latency);
241+
loop_read = hdfsRead(_handle->fs(), _handle->file(), to + has_read,
242+
static_cast<int32_t>(bytes_req - has_read));
243+
}
235244
if (loop_read < 0) {
236245
// invoker maybe just skip Status.NotFound and continue
237246
// so we need distinguish between it and other kinds of errors

be/src/io/fs/hdfs_file_system.cpp

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@
4040
#include "io/hdfs_builder.h"
4141
#include "io/hdfs_util.h"
4242
#include "runtime/exec_env.h"
43+
#include "util/bvar_helper.h"
4344
#include "util/obj_lru_cache.h"
4445
#include "util/slice.h"
4546

@@ -120,7 +121,11 @@ Status HdfsFileSystem::open_file_internal(const Path& file, FileReaderSPtr* read
120121
Status HdfsFileSystem::create_directory_impl(const Path& dir, bool failed_if_exists) {
121122
CHECK_HDFS_HANDLER(_fs_handler);
122123
Path real_path = convert_path(dir, _fs_name);
123-
int res = hdfsCreateDirectory(_fs_handler->hdfs_fs, real_path.string().c_str());
124+
int res;
125+
{
126+
SCOPED_BVAR_LATENCY(hdfs_bvar::hdfs_create_dir_latency);
127+
res = hdfsCreateDirectory(_fs_handler->hdfs_fs, real_path.string().c_str());
128+
}
124129
if (res == -1) {
125130
return Status::IOError("failed to create directory {}: {}", dir.native(), hdfs_error());
126131
}

be/test/format_v2/table/iceberg_reader_test.cpp

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1203,14 +1203,14 @@ TIcebergDeleteFileDesc make_iceberg_position_delete_file(const std::string& path
12031203

12041204
TIcebergDeleteFileDesc make_iceberg_equality_delete_file(
12051205
const std::string& path, const std::vector<int32_t>& field_ids,
1206-
TFileFormatType::type file_format = TFileFormatType::FORMAT_PARQUET) {
1206+
TFileFormatType::type file_format = TFileFormatType::FORMAT_PARQUET,
1207+
int64_t file_size = -1) {
12071208
TIcebergDeleteFileDesc delete_file;
12081209
delete_file.__set_content(2);
12091210
delete_file.__set_path(path);
12101211
delete_file.__set_field_ids(field_ids);
12111212
delete_file.__set_file_format(file_format);
1212-
// Set file_size to actual file size, simulating FE propagation from iceberg manifest
1213-
delete_file.__set_file_size(static_cast<int64_t>(std::filesystem::file_size(path)));
1213+
delete_file.__set_file_size(file_size);
12141214
return delete_file;
12151215
}
12161216

@@ -5062,7 +5062,9 @@ TEST(IcebergV2ReaderTest, IcebergEqualityDeleteFileSizePropagatedToReader) {
50625062
auto split_options = build_split_options(file_path);
50635063
split_options.cache = &cache;
50645064
split_options.current_range.__set_table_format_params(make_iceberg_table_format_desc(
5065-
file_path, {make_iceberg_equality_delete_file(delete_file_path, {0})}));
5065+
file_path, {make_iceberg_equality_delete_file(delete_file_path, {0},
5066+
TFileFormatType::FORMAT_PARQUET,
5067+
delete_file_size)}));
50665068
ASSERT_TRUE(reader.prepare_split(split_options).ok());
50675069

50685070
Block block = build_table_block(projected_columns);

be/test/io/fs/file_handle_cache_test.cpp

Lines changed: 0 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -55,17 +55,4 @@ TEST(FileHandleCacheTest, InitWithKnownFileSizeDoesNotOpenFile) {
5555
EXPECT_EQ(handle.file(), nullptr);
5656
}
5757

58-
// Verify that build_iceberg_delete_file_range treats file_size <= 0 as unknown (-1).
59-
// This covers the case where thrift optional file_size defaults to 0.
60-
TEST(FileHandleCacheTest, BuildDeleteFileRangeTreatsZeroAsUnknown) {
61-
auto range_known = build_iceberg_delete_file_range("s3://b/f.parquet", 1024);
62-
EXPECT_EQ(range_known.file_size, 1024);
63-
64-
auto range_zero = build_iceberg_delete_file_range("s3://b/f.parquet", 0);
65-
EXPECT_EQ(range_zero.file_size, -1);
66-
67-
auto range_neg = build_iceberg_delete_file_range("s3://b/f.parquet", -1);
68-
EXPECT_EQ(range_neg.file_size, -1);
69-
}
70-
7158
} // namespace doris::io

0 commit comments

Comments
 (0)