This is an automated email from the ASF dual-hosted git repository.
Gabriel39 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 dfa3b018a4e [fix](iceberg) make iceberg deletion vector stable.
(#65676)
dfa3b018a4e is described below
commit dfa3b018a4e3ee6cc99f7907d1d2ecc376d919c4
Author: daidai <[email protected]>
AuthorDate: Mon Jul 20 11:46:25 2026 +0800
[fix](iceberg) make iceberg deletion vector stable. (#65676)
### What problem does this PR solve?
Related PR: #59272
Doris already supports Iceberg v3 deletion vectors, but several
stability gaps remain:
- Invalid Iceberg/Paimon offsets and lengths could reach cache lookup or
memory allocation.
- Iceberg Puffin deletion-vector CRC32 was not verified.
- Validation was inconsistent across legacy readers, format-v2 readers,
normal scans, and the `position_deletes` path.
- The Iceberg reader enforced a 1 GiB limit, while the writer did not.
### What is changed?
- Validate Iceberg/Paimon deletion-vector descriptors before cache
lookup and allocation.
- Verify Iceberg Puffin CRC32 before bitmap decoding.
- Report malformed Paimon framing as data-quality errors.
- Validate Puffin metadata in normal FE scans and the `position_deletes`
path.
- Align the Iceberg writer/reader 1 GiB limit and separate
Iceberg/Paimon limits.
- Add BE, FE, and format-v2 boundary tests.
### Release note
Harden Iceberg and Paimon deletion-vector validation and consistently
enforce the 1 GiB Iceberg deletion-vector limit on write and read paths.
### Check List (For Author)
- Test <!-- At least one of them must be included. -->
- [ ] Regression test
- [x] Unit Test
- [ ] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change.
- [ ] No code files have been changed.
- [ ] Other reason <!-- Add your reason? -->
- Behavior changed:
- [ ] No.
- [x] Yes. <!-- Explain the behavior change -->
- Does this need documentation?
- [x] No.
- [ ] Yes. <!-- Add document PR link here. eg:
https://github.com/apache/doris-website/pull/1214 -->
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label <!-- Add branch pick label that this PR
should merge into -->
---
be/src/exec/sink/viceberg_delete_sink.cpp | 20 +++-
be/src/exec/sink/viceberg_delete_sink.h | 4 +
be/src/format/table/deletion_vector.h | 12 +++
be/src/format/table/deletion_vector_reader.cpp | 82 +++++++++++++++-
be/src/format/table/deletion_vector_reader.h | 8 +-
.../table/iceberg_delete_file_reader_helper.cpp | 47 ++++++---
.../table/iceberg_delete_file_reader_helper.h | 3 +
be/src/format/table/iceberg_reader_mixin.h | 3 +
be/src/format/table/paimon_reader.cpp | 41 +++++---
be/src/format/table/paimon_reader.h | 4 +
be/src/format_v2/table/iceberg_reader.cpp | 8 +-
be/src/format_v2/table/paimon_reader.cpp | 4 +-
be/test/exec/sink/viceberg_delete_sink_test.cpp | 14 +++
.../iceberg_delete_file_reader_helper_test.cpp | 109 +++++++++++++++++++--
be/test/format/table/paimon_cpp_reader_test.cpp | 52 +++++++++-
be/test/format_v2/table/iceberg_reader_test.cpp | 25 ++++-
be/test/format_v2/table/paimon_reader_test.cpp | 14 +++
.../iceberg/source/IcebergDeleteFileFilter.java | 32 +++++-
.../datasource/iceberg/source/IcebergScanNode.java | 8 +-
.../source/IcebergDeleteFileFilterTest.java | 48 +++++++++
.../iceberg/source/IcebergScanNodeTest.java | 30 ++++++
21 files changed, 520 insertions(+), 48 deletions(-)
diff --git a/be/src/exec/sink/viceberg_delete_sink.cpp
b/be/src/exec/sink/viceberg_delete_sink.cpp
index e2b142db8d9..ef167ecdbca 100644
--- a/be/src/exec/sink/viceberg_delete_sink.cpp
+++ b/be/src/exec/sink/viceberg_delete_sink.cpp
@@ -35,6 +35,7 @@
#include "core/data_type/data_type_struct.h"
#include "exec/common/endian.h"
#include "exprs/vexpr.h"
+#include "format/table/deletion_vector.h"
#include "format/table/iceberg_delete_file_reader_helper.h"
#include "format/transformer/vfile_format_transformer.h"
#include "io/file_factory.h"
@@ -106,6 +107,20 @@ Status load_rewritable_delete_rows(RuntimeState* state,
RuntimeProfile* profile,
} // namespace
+Status calculate_iceberg_deletion_vector_content_size(size_t bitmap_size,
int64_t* content_size) {
+ DORIS_CHECK(content_size != nullptr);
+ constexpr size_t max_bitmap_size =
static_cast<size_t>(MAX_ICEBERG_DELETION_VECTOR_BYTES) -
+
ICEBERG_DELETION_VECTOR_BLOB_OVERHEAD_BYTES;
+ if (bitmap_size > max_bitmap_size) {
+ return Status::NotSupported(
+ "Iceberg deletion vector bitmap size exceeds Doris supported
limit: {}, "
+ "maximum bitmap size: {}, content size limit: {}",
+ bitmap_size, max_bitmap_size,
MAX_ICEBERG_DELETION_VECTOR_BYTES);
+ }
+ *content_size = static_cast<int64_t>(bitmap_size +
ICEBERG_DELETION_VECTOR_BLOB_OVERHEAD_BYTES);
+ return Status::OK();
+}
+
VIcebergDeleteSink::VIcebergDeleteSink(const TDataSink& t_sink,
const VExprContextSPtrs& output_exprs,
std::shared_ptr<Dependency> dep,
@@ -592,7 +607,8 @@ Status VIcebergDeleteSink::_write_deletion_vector_files(
blob.partition_spec_id = deletion.partition_spec_id;
blob.partition_data_json = deletion.partition_data_json;
blob.merged_count = static_cast<int64_t>(merged_rows.cardinality());
- blob.content_size_in_bytes = static_cast<int64_t>(4 + 4 + bitmap_size
+ 4);
+ RETURN_IF_ERROR(calculate_iceberg_deletion_vector_content_size(
+ bitmap_size, &blob.content_size_in_bytes));
blob.blob_data.resize(static_cast<size_t>(blob.content_size_in_bytes));
merged_rows.write(blob.blob_data.data() + 8);
@@ -604,7 +620,7 @@ Status VIcebergDeleteSink::_write_deletion_vector_files(
uint32_t crc = static_cast<uint32_t>(
::crc32(0, reinterpret_cast<const
Bytef*>(blob.blob_data.data() + 4),
- 4 + (uInt)bitmap_size));
+ static_cast<uInt>(4 + bitmap_size)));
BigEndian::Store32(blob.blob_data.data() + 8 + bitmap_size, crc);
blobs.emplace_back(std::move(blob));
}
diff --git a/be/src/exec/sink/viceberg_delete_sink.h
b/be/src/exec/sink/viceberg_delete_sink.h
index 22ae98cc288..647ee92f525 100644
--- a/be/src/exec/sink/viceberg_delete_sink.h
+++ b/be/src/exec/sink/viceberg_delete_sink.h
@@ -20,6 +20,8 @@
#include <gen_cpp/DataSinks_types.h>
#include <gen_cpp/PlanNodes_types.h>
+#include <cstddef>
+#include <cstdint>
#include <map>
#include <memory>
#include <string>
@@ -41,6 +43,8 @@ class Dependency;
class VIcebergDeleteFileWriter;
+Status calculate_iceberg_deletion_vector_content_size(size_t bitmap_size,
int64_t* content_size);
+
struct IcebergFileDeletion {
IcebergFileDeletion() = default;
IcebergFileDeletion(int32_t spec_id, std::string data_json)
diff --git a/be/src/format/table/deletion_vector.h
b/be/src/format/table/deletion_vector.h
index 2c89771b2e6..12ec7dbc6eb 100644
--- a/be/src/format/table/deletion_vector.h
+++ b/be/src/format/table/deletion_vector.h
@@ -17,10 +17,22 @@
#pragma once
+#include <cstddef>
+#include <cstdint>
+
#include "roaring/roaring64map.hh"
namespace doris {
+// Doris materializes an Iceberg deletion vector as one buffer. This is an
implementation limit,
+// not an Iceberg format limit.
+inline constexpr int64_t MAX_ICEBERG_DELETION_VECTOR_BYTES = 1L << 30;
+inline constexpr size_t ICEBERG_DELETION_VECTOR_BLOB_OVERHEAD_BYTES = 12;
+
+// Paimon v1 uses a run-optimized 32-bit Roaring bitmap whose maximum
serialized size is below this
+// limit. Keep its guard independent from the Iceberg capability limit.
+inline constexpr int64_t MAX_PAIMON_DELETION_VECTOR_BYTES = 1L << 30;
+
// A deletion vector is already a bitmap on the wire. Keep decoded DVs
compressed in the
// query-local cache instead of expanding every set bit into an int64_t.
Position delete files use
// a different representation because their input is a stream of (file_path,
row_position) rows.
diff --git a/be/src/format/table/deletion_vector_reader.cpp
b/be/src/format/table/deletion_vector_reader.cpp
index 65ce2831a09..8e8404839be 100644
--- a/be/src/format/table/deletion_vector_reader.cpp
+++ b/be/src/format/table/deletion_vector_reader.cpp
@@ -17,11 +17,77 @@
#include "format/table/deletion_vector_reader.h"
+#include <limits>
+
#include "rapidjson/document.h"
#include "rapidjson/stringbuffer.h"
#include "util/block_compression.h"
namespace doris {
+namespace {
+
+constexpr int64_t ICEBERG_DELETION_VECTOR_MIN_BYTES =
+ static_cast<int64_t>(ICEBERG_DELETION_VECTOR_BLOB_OVERHEAD_BYTES);
+constexpr int64_t PAIMON_DELETION_VECTOR_MIN_BYTES = 8;
+constexpr int64_t PAIMON_LENGTH_PREFIX_BYTES = 4;
+
+enum class DeletionVectorSizeLimitStatus {
+ DATA_QUALITY,
+ NOT_SUPPORTED,
+};
+
+Status validate_deletion_vector_read_range(const char* description, int64_t
offset, int64_t size,
+ int64_t min_size, int64_t max_size,
+ DeletionVectorSizeLimitStatus
size_limit_status,
+ size_t& bytes_read) {
+ if (offset < 0) {
+ return Status::DataQualityError("{} offset must be non-negative: {}",
description, offset);
+ }
+ if (size < min_size) {
+ return Status::DataQualityError("{} size too small: {}, minimum: {}",
description, size,
+ min_size);
+ }
+ if (size > max_size) {
+ if (size_limit_status == DeletionVectorSizeLimitStatus::NOT_SUPPORTED)
{
+ return Status::NotSupported("{} size exceeds Doris supported
limit: {}, limit: {}",
+ description, size, max_size);
+ }
+ return Status::DataQualityError("{} size exceeds limit: {}, limit:
{}", description, size,
+ max_size);
+ }
+ if (offset > std::numeric_limits<int64_t>::max() - size) {
+ return Status::DataQualityError("{} offset plus size overflows: offset
{}, size {}",
+ description, offset, size);
+ }
+ bytes_read = static_cast<size_t>(size);
+ return Status::OK();
+}
+
+} // namespace
+
+Status validate_iceberg_deletion_vector_read_range(int64_t offset, int64_t
size,
+ size_t& bytes_read) {
+ return validate_deletion_vector_read_range(
+ "Iceberg deletion vector", offset, size,
ICEBERG_DELETION_VECTOR_MIN_BYTES,
+ MAX_ICEBERG_DELETION_VECTOR_BYTES,
DeletionVectorSizeLimitStatus::NOT_SUPPORTED,
+ bytes_read);
+}
+
+Status validate_paimon_deletion_vector_read_range(int64_t offset, int64_t
length,
+ size_t& bytes_read) {
+ if (length < 0) {
+ return Status::DataQualityError("Paimon deletion vector length must be
non-negative: {}",
+ length);
+ }
+ if (length > std::numeric_limits<int64_t>::max() -
PAIMON_LENGTH_PREFIX_BYTES) {
+ return Status::DataQualityError("Paimon deletion vector length
overflows: {}", length);
+ }
+ return validate_deletion_vector_read_range(
+ "Paimon deletion vector", offset, length +
PAIMON_LENGTH_PREFIX_BYTES,
+ PAIMON_DELETION_VECTOR_MIN_BYTES, MAX_PAIMON_DELETION_VECTOR_BYTES,
+ DeletionVectorSizeLimitStatus::DATA_QUALITY, bytes_read);
+}
+
DeletionVectorReader::~DeletionVectorReader() {
// The file reader may retain the child IOContext. Destroy it before
merging and before the
// child statistics storage goes away.
@@ -60,11 +126,25 @@ Status DeletionVectorReader::open() {
return Status::OK();
}
+ if (_desc.start_offset < 0 || _desc.size < 0) {
+ return Status::DataQualityError(
+ "Deletion vector range must be non-negative: path {}, offset
{}, size {}",
+ _desc.path, _desc.start_offset, _desc.size);
+ }
+
_init_system_properties();
_init_file_description();
RETURN_IF_ERROR(_create_file_reader());
- _file_size = _file_reader->size();
+ const size_t file_size = _file_reader->size();
+ const size_t start_offset = static_cast<size_t>(_desc.start_offset);
+ const size_t range_size = static_cast<size_t>(_desc.size);
+ if (start_offset > file_size || range_size > file_size - start_offset) {
+ return Status::DataQualityError(
+ "Deletion vector range exceeds file size: path {}, offset {},
size {}, file size "
+ "{}",
+ _desc.path, _desc.start_offset, _desc.size, file_size);
+ }
_is_opened = true;
return Status::OK();
}
diff --git a/be/src/format/table/deletion_vector_reader.h
b/be/src/format/table/deletion_vector_reader.h
index 55f0099b6fe..c187e287168 100644
--- a/be/src/format/table/deletion_vector_reader.h
+++ b/be/src/format/table/deletion_vector_reader.h
@@ -17,6 +17,7 @@
#pragma once
+#include <cstddef>
#include <cstdint>
#include <memory>
#include <string>
@@ -37,6 +38,12 @@ struct IOContext;
} // namespace io
namespace doris {
+Status validate_iceberg_deletion_vector_read_range(int64_t offset, int64_t
size,
+ size_t& bytes_read);
+
+Status validate_paimon_deletion_vector_read_range(int64_t offset, int64_t
length,
+ size_t& bytes_read);
+
struct DeleteFileDesc {
enum class Format {
PAIMON,
@@ -111,7 +118,6 @@ private:
io::FileSystemProperties _system_properties;
io::FileDescription _file_description;
io::FileReaderSPtr _file_reader;
- int64_t _file_size = 0;
bool _is_opened = false;
};
} // namespace doris
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 9f4c0482e2a..f8828e402b0 100644
--- a/be/src/format/table/iceberg_delete_file_reader_helper.cpp
+++ b/be/src/format/table/iceberg_delete_file_reader_helper.cpp
@@ -48,6 +48,7 @@
#include "runtime/runtime_state.h"
#include "storage/predicate/column_predicate.h"
#include "util/debug_points.h"
+#include "util/hash_util.hpp"
namespace doris {
@@ -55,6 +56,7 @@ namespace {
constexpr const char* ICEBERG_FILE_PATH = "file_path";
constexpr const char* ICEBERG_ROW_POS = "pos";
+constexpr size_t ICEBERG_DELETION_VECTOR_MIN_BYTES = 12;
const std::vector<std::string> DELETE_COL_NAMES {ICEBERG_FILE_PATH,
ICEBERG_ROW_POS};
std::unordered_map<std::string, uint32_t> DELETE_COL_NAME_TO_BLOCK_IDX =
{{ICEBERG_FILE_PATH, 0},
@@ -177,14 +179,14 @@ Status decode_deletion_vector_buffer(const char* buf,
size_t buffer_size,
if (buf == nullptr || rows_to_delete == nullptr) {
return Status::InvalidArgument("invalid deletion vector decode
arguments");
}
- if (buffer_size < 12) {
+ if (buffer_size < ICEBERG_DELETION_VECTOR_MIN_BYTES) {
return Status::DataQualityError("Deletion vector file size too small:
{}", buffer_size);
}
- auto total_length = BigEndian::Load32(buf);
- if (total_length + 8 != buffer_size) {
+ const uint32_t total_length = BigEndian::Load32(buf);
+ if (static_cast<uint64_t>(total_length) + 8 != buffer_size) {
return Status::DataQualityError("Deletion vector length mismatch,
expected: {}, actual: {}",
- total_length + 8, buffer_size);
+ static_cast<uint64_t>(total_length) +
8, buffer_size);
}
constexpr static char MAGIC_NUMBER[] = {'\xD1', '\xD3', '\x39', '\x64'};
@@ -192,8 +194,17 @@ Status decode_deletion_vector_buffer(const char* buf,
size_t buffer_size,
return Status::DataQualityError("Deletion vector magic number
mismatch");
}
+ const uint32_t expected_crc = BigEndian::Load32(buf + sizeof(total_length)
+ total_length);
+ const uint32_t actual_crc =
+ HashUtil::zlib_crc_hash(buf + sizeof(total_length), total_length,
0);
+ if (actual_crc != expected_crc) {
+ return Status::DataQualityError("Deletion vector CRC32 mismatch,
expected: {}, actual: {}",
+ expected_crc, actual_crc);
+ }
+
try {
- *rows_to_delete |= roaring::Roaring64Map::readSafe(buf + 8,
buffer_size - 12);
+ *rows_to_delete |= roaring::Roaring64Map::readSafe(
+ buf + 8, buffer_size - ICEBERG_DELETION_VECTOR_MIN_BYTES);
} catch (const std::runtime_error& e) {
return Status::DataQualityError("Decode roaring bitmap failed, {}",
e.what());
}
@@ -245,6 +256,18 @@ std::string build_iceberg_deletion_vector_cache_key(const
std::string& data_file
delete_file.content_size_in_bytes);
}
+Status validate_iceberg_deletion_vector_descriptor(const
TIcebergDeleteFileDesc& delete_file,
+ size_t& bytes_read) {
+ if (!delete_file.__isset.path || !delete_file.__isset.content_offset ||
+ !delete_file.__isset.content_size_in_bytes) {
+ return Status::DataQualityError(
+ "Iceberg deletion vector descriptor misses "
+ "path/content_offset/content_size_in_bytes");
+ }
+ return validate_iceberg_deletion_vector_read_range(
+ delete_file.content_offset, delete_file.content_size_in_bytes,
bytes_read);
+}
+
Status read_iceberg_position_delete_file(const TIcebergDeleteFileDesc&
delete_file,
const IcebergDeleteFileReaderOptions&
options,
IcebergPositionDeleteVisitor*
visitor) {
@@ -318,9 +341,8 @@ Status read_iceberg_deletion_vector(const
TIcebergDeleteFileDesc& delete_file,
options.io_ctx == nullptr || rows_to_delete == nullptr) {
return Status::InvalidArgument("invalid deletion vector reader
options");
}
- if (!delete_file.__isset.content_offset ||
!delete_file.__isset.content_size_in_bytes) {
- return Status::InternalError("Deletion vector is missing content
offset or length");
- }
+ size_t bytes_read = 0;
+ RETURN_IF_ERROR(validate_iceberg_deletion_vector_descriptor(delete_file,
bytes_read));
DBUG_EXECUTE_IF("IcebergDeleteFileReader.read_deletion_vector.io_error",
{ return Status::IOError("injected Iceberg deletion vector
read failure"); });
DBUG_EXECUTE_IF("IcebergDeleteFileReader.read_deletion_vector.should_stop",
@@ -333,18 +355,19 @@ Status read_iceberg_deletion_vector(const
TIcebergDeleteFileDesc& delete_file,
delete_range.start_offset = delete_file.content_offset;
delete_range.size = delete_file.content_size_in_bytes;
+ // Iceberg v3 deletion-vector-v1 blobs are uncompressed and metadata
provides the exact range.
+ // Parse the Puffin footer first if future blob types or compression
codecs are supported.
DeletionVectorReader dv_reader(options.state, options.profile,
*options.scan_params,
delete_range, options.io_ctx);
RETURN_IF_ERROR(dv_reader.open());
- std::vector<char> buf(delete_range.size);
- const auto read_status = dv_reader.read_at(delete_range.start_offset,
- {buf.data(),
cast_set<size_t>(delete_range.size)});
+ std::vector<char> buf(bytes_read);
+ const auto read_status = dv_reader.read_at(delete_range.start_offset,
{buf.data(), bytes_read});
if (options.deletion_vector_file_cache_stats != nullptr) {
options.deletion_vector_file_cache_stats->merge_from(dv_reader.file_cache_statistics());
}
RETURN_IF_ERROR(read_status);
- return decode_deletion_vector_buffer(buf.data(), delete_range.size,
rows_to_delete);
+ return decode_deletion_vector_buffer(buf.data(), bytes_read,
rows_to_delete);
}
Status decode_iceberg_deletion_vector_buffer(const char* buf, size_t
buffer_size,
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 0b4f601c0e8..adc0ef196f4 100644
--- a/be/src/format/table/iceberg_delete_file_reader_helper.h
+++ b/be/src/format/table/iceberg_delete_file_reader_helper.h
@@ -74,6 +74,9 @@ bool is_iceberg_deletion_vector(const TIcebergDeleteFileDesc&
delete_file);
std::string build_iceberg_deletion_vector_cache_key(const std::string&
data_file_path,
const
TIcebergDeleteFileDesc& delete_file);
+Status validate_iceberg_deletion_vector_descriptor(const
TIcebergDeleteFileDesc& delete_file,
+ size_t& bytes_read);
+
Status decode_iceberg_deletion_vector_buffer(const char* buf, size_t
buffer_size,
DeletionVector* rows_to_delete);
diff --git a/be/src/format/table/iceberg_reader_mixin.h
b/be/src/format/table/iceberg_reader_mixin.h
index aee211aee01..b9d98ebb660 100644
--- a/be/src/format/table/iceberg_reader_mixin.h
+++ b/be/src/format/table/iceberg_reader_mixin.h
@@ -996,6 +996,9 @@ Status
IcebergReaderMixin<BaseReader>::_gen_position_delete_file_range(
template <typename BaseReader>
Status IcebergReaderMixin<BaseReader>::read_deletion_vector(
const std::string& data_file_path, const TIcebergDeleteFileDesc&
delete_file_desc) {
+ size_t bytes_read = 0;
+
RETURN_IF_ERROR(validate_iceberg_deletion_vector_descriptor(delete_file_desc,
bytes_read));
+
Status create_status = Status::OK();
SCOPED_TIMER(_iceberg_profile.delete_files_read_time);
bool decoded_cache_hit = false;
diff --git a/be/src/format/table/paimon_reader.cpp
b/be/src/format/table/paimon_reader.cpp
index 4953078c32d..a73d65f8336 100644
--- a/be/src/format/table/paimon_reader.cpp
+++ b/be/src/format/table/paimon_reader.cpp
@@ -42,26 +42,41 @@ std::string build_paimon_deletion_vector_cache_key(const
TPaimonDeletionFileDesc
deletion_file.length);
}
+Status validate_paimon_deletion_vector_descriptor(const
TPaimonDeletionFileDesc& deletion_file,
+ size_t& bytes_read) {
+ if (!deletion_file.__isset.path || !deletion_file.__isset.offset ||
+ !deletion_file.__isset.length) {
+ return Status::DataQualityError(
+ "Paimon deletion file descriptor misses path/offset/length");
+ }
+ return validate_paimon_deletion_vector_read_range(deletion_file.offset,
deletion_file.length,
+ bytes_read);
+}
+
Status decode_paimon_deletion_vector_buffer(const char* buf, size_t
buffer_size,
DeletionVector* deletion_vector) {
if (deletion_vector == nullptr) {
- return Status::InvalidArgument("deletion_vector must not be null");
+ return Status::InvalidArgument("deletion vector output must not be
null");
+ }
+ if (buf == nullptr) {
+ return Status::DataQualityError("Paimon deletion vector blob is null");
}
if (buffer_size < 8) [[unlikely]] {
return Status::DataQualityError("Deletion vector file size too small:
{}", buffer_size);
}
const uint32_t actual_length = BigEndian::Load32(buf);
- if (actual_length + 4 != buffer_size) [[unlikely]] {
- return Status::RuntimeError(
- "DeletionVector deserialize error: length not match, "
- "actual length: {}, expect length: {}",
- actual_length, buffer_size - 4);
+ if (static_cast<uint64_t>(actual_length) + 4 != buffer_size) [[unlikely]] {
+ return Status::DataQualityError(
+ "Paimon deletion vector length mismatch, expected: {}, actual:
{}",
+ static_cast<uint64_t>(actual_length) + 4, buffer_size);
}
if (memcmp(buf + sizeof(actual_length), PAIMON_BITMAP_MAGIC, 4) != 0)
[[unlikely]] {
- return Status::RuntimeError("DeletionVector deserialize error: invalid
magic number {}",
- BigEndian::Load32(buf +
sizeof(actual_length)));
+ return Status::DataQualityError(
+ "Paimon deletion vector magic number mismatch, expected: {},
actual: {}",
+ BigEndian::Load32(PAIMON_BITMAP_MAGIC),
+ BigEndian::Load32(buf + sizeof(actual_length)));
}
roaring::Roaring roaring_bitmap;
@@ -158,6 +173,8 @@ Status PaimonOrcReader::_init_deletion_vector() {
set_push_down_agg_type(TPushAggOp::NONE);
}
const auto& deletion_file = table_desc.deletion_file;
+ size_t bytes_read = 0;
+ RETURN_IF_ERROR(validate_paimon_deletion_vector_descriptor(deletion_file,
bytes_read));
Status create_status = Status::OK();
@@ -172,7 +189,7 @@ Status PaimonOrcReader::_init_deletion_vector() {
delete_range.__set_fs_name(get_scan_range().fs_name);
delete_range.path = deletion_file.path;
delete_range.start_offset = deletion_file.offset;
- delete_range.size = deletion_file.length + 4;
+ delete_range.size = static_cast<int64_t>(bytes_read);
delete_range.file_size = -1;
DeletionVectorReader dv_reader(get_state(), get_profile(),
get_scan_params(),
@@ -182,7 +199,6 @@ Status PaimonOrcReader::_init_deletion_vector() {
return nullptr;
}
- size_t bytes_read = deletion_file.length + 4;
std::vector<char> buffer(bytes_read);
create_status =
dv_reader.read_at(deletion_file.offset,
{buffer.data(), bytes_read});
@@ -264,6 +280,8 @@ Status PaimonParquetReader::_init_deletion_vector() {
set_push_down_agg_type(TPushAggOp::NONE);
}
const auto& deletion_file = table_desc.deletion_file;
+ size_t bytes_read = 0;
+ RETURN_IF_ERROR(validate_paimon_deletion_vector_descriptor(deletion_file,
bytes_read));
Status create_status = Status::OK();
@@ -278,7 +296,7 @@ Status PaimonParquetReader::_init_deletion_vector() {
delete_range.__set_fs_name(get_scan_range().fs_name);
delete_range.path = deletion_file.path;
delete_range.start_offset = deletion_file.offset;
- delete_range.size = deletion_file.length + 4;
+ delete_range.size = static_cast<int64_t>(bytes_read);
delete_range.file_size = -1;
DeletionVectorReader dv_reader(get_state(), get_profile(),
get_scan_params(),
@@ -288,7 +306,6 @@ Status PaimonParquetReader::_init_deletion_vector() {
return nullptr;
}
- size_t bytes_read = deletion_file.length + 4;
std::vector<char> buffer(bytes_read);
create_status =
dv_reader.read_at(deletion_file.offset,
{buffer.data(), bytes_read});
diff --git a/be/src/format/table/paimon_reader.h
b/be/src/format/table/paimon_reader.h
index 0f11825ce89..746ee1adec1 100644
--- a/be/src/format/table/paimon_reader.h
+++ b/be/src/format/table/paimon_reader.h
@@ -19,6 +19,7 @@
#include <gen_cpp/PlanNodes_types.h>
+#include <cstddef>
#include <memory>
#include <string>
#include <utility>
@@ -35,6 +36,9 @@ class ShardedKVCache;
std::string build_paimon_deletion_vector_cache_key(const
TPaimonDeletionFileDesc& deletion_file);
+Status validate_paimon_deletion_vector_descriptor(const
TPaimonDeletionFileDesc& deletion_file,
+ size_t& bytes_read);
+
Status decode_paimon_deletion_vector_buffer(const char* buf, size_t
buffer_size,
DeletionVector* deletion_vector);
diff --git a/be/src/format_v2/table/iceberg_reader.cpp
b/be/src/format_v2/table/iceberg_reader.cpp
index 75df8e4bd6c..b6580a12dfc 100644
--- a/be/src/format_v2/table/iceberg_reader.cpp
+++ b/be/src/format_v2/table/iceberg_reader.cpp
@@ -387,10 +387,8 @@ Status
IcebergTableReader::_parse_deletion_vector_file(const TTableFormatFileDes
if (deletion_vector == nullptr) {
return Status::OK();
}
- if (!deletion_vector->__isset.content_offset ||
- !deletion_vector->__isset.content_size_in_bytes) {
- return Status::InternalError("Deletion vector is missing content
offset or length");
- }
+ size_t bytes_read = 0;
+
RETURN_IF_ERROR(validate_iceberg_deletion_vector_descriptor(*deletion_vector,
bytes_read));
const std::string data_file_path =
iceberg_params.__isset.original_file_path
?
iceberg_params.original_file_path
@@ -398,7 +396,7 @@ Status
IcebergTableReader::_parse_deletion_vector_file(const TTableFormatFileDes
desc->key = build_iceberg_deletion_vector_cache_key(data_file_path,
*deletion_vector);
desc->path = deletion_vector->path;
desc->start_offset = deletion_vector->content_offset;
- desc->size = deletion_vector->content_size_in_bytes;
+ desc->size = static_cast<int64_t>(bytes_read);
desc->file_size = -1;
desc->format = DeleteFileDesc::Format::ICEBERG;
*has_delete_file = true;
diff --git a/be/src/format_v2/table/paimon_reader.cpp
b/be/src/format_v2/table/paimon_reader.cpp
index 74a38515622..0b3e410b87c 100644
--- a/be/src/format_v2/table/paimon_reader.cpp
+++ b/be/src/format_v2/table/paimon_reader.cpp
@@ -75,11 +75,13 @@ Status PaimonReader::_parse_deletion_vector_file(const
TTableFormatFileDesc& t_d
return Status::OK();
}
const auto& deletion_file = table_desc.deletion_file;
+ size_t bytes_read = 0;
+ RETURN_IF_ERROR(validate_paimon_deletion_vector_descriptor(deletion_file,
bytes_read));
desc->key = build_paimon_deletion_vector_cache_key(deletion_file);
desc->path = deletion_file.path;
desc->start_offset = deletion_file.offset;
- desc->size = deletion_file.length + 4;
+ desc->size = static_cast<int64_t>(bytes_read);
desc->file_size = -1;
desc->format = DeleteFileDesc::Format::PAIMON;
*has_delete_file = true;
diff --git a/be/test/exec/sink/viceberg_delete_sink_test.cpp
b/be/test/exec/sink/viceberg_delete_sink_test.cpp
index d9fc5086503..7faa77ed702 100644
--- a/be/test/exec/sink/viceberg_delete_sink_test.cpp
+++ b/be/test/exec/sink/viceberg_delete_sink_test.cpp
@@ -33,6 +33,7 @@
#include "core/data_type/data_type_string.h"
#include "core/data_type/data_type_struct.h"
#include "exec/common/endian.h"
+#include "format/table/deletion_vector.h"
#include "gen_cpp/DataSinks_types.h"
#include "gen_cpp/Types_types.h"
#include "runtime/runtime_state.h"
@@ -507,6 +508,19 @@ TEST_F(VIcebergDeleteSinkTest, TestUnsupportedDeleteType) {
ASSERT_FALSE(status.ok());
}
+TEST_F(VIcebergDeleteSinkTest, TestValidateDeletionVectorContentSize) {
+ constexpr size_t max_bitmap_size =
static_cast<size_t>(MAX_ICEBERG_DELETION_VECTOR_BYTES) -
+
ICEBERG_DELETION_VECTOR_BLOB_OVERHEAD_BYTES;
+ int64_t content_size = 0;
+ ASSERT_TRUE(
+ calculate_iceberg_deletion_vector_content_size(max_bitmap_size,
&content_size).ok());
+ ASSERT_EQ(MAX_ICEBERG_DELETION_VECTOR_BYTES, content_size);
+
+ const auto unsupported_status =
+ calculate_iceberg_deletion_vector_content_size(max_bitmap_size +
1, &content_size);
+ ASSERT_TRUE(unsupported_status.is<ErrorCode::NOT_IMPLEMENTED_ERROR>()) <<
unsupported_status;
+}
+
TEST_F(VIcebergDeleteSinkTest, TestWriteDeletionVectorsToSingleSharedPuffin) {
std::filesystem::path temp_dir = std::filesystem::temp_directory_path() /
("iceberg_delete_sink_test_" +
generate_uuid_string());
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 7c517d49f61..841e01f9cc9 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
@@ -27,6 +27,7 @@
#include <cstring>
#include <filesystem>
#include <fstream>
+#include <limits>
#include <optional>
#include <string>
#include <thread>
@@ -36,12 +37,14 @@
#include "cloud/config.h"
#include "common/config.h"
#include "exec/common/endian.h"
+#include "format/table/deletion_vector_reader.h"
#include "io/fs/file_meta_cache.h"
#include "roaring/roaring64map.hh"
#include "runtime/runtime_profile.h"
#include "runtime/runtime_state.h"
#include "runtime/thread_context.h"
#include "testutil/mock/mock_runtime_state.h"
+#include "util/hash_util.hpp"
namespace doris {
@@ -112,8 +115,8 @@ TIcebergDeleteFileDesc make_iceberg_deletion_vector(const
std::string& path, int
return delete_file;
}
-int64_t write_iceberg_deletion_vector_file(const std::string& file_path,
- const std::vector<uint64_t>&
deleted_positions) {
+std::vector<char> build_iceberg_deletion_vector_blob(
+ const std::vector<uint64_t>& deleted_positions) {
roaring::Roaring64Map rows;
for (const auto position : deleted_positions) {
rows.add(position);
@@ -127,7 +130,14 @@ int64_t write_iceberg_deletion_vector_file(const
std::string& file_path,
BigEndian::Store32(blob.data(), total_length);
constexpr char DV_MAGIC[] = {'\xD1', '\xD3', '\x39', '\x64'};
memcpy(blob.data() + 4, DV_MAGIC, 4);
- BigEndian::Store32(blob.data() + 8 + bitmap_size, 0);
+ const uint32_t crc = HashUtil::zlib_crc_hash(blob.data() + 4,
total_length, 0);
+ BigEndian::Store32(blob.data() + 8 + bitmap_size, crc);
+ return blob;
+}
+
+int64_t write_iceberg_deletion_vector_file(const std::string& file_path,
+ const std::vector<uint64_t>&
deleted_positions) {
+ const auto blob = build_iceberg_deletion_vector_blob(deleted_positions);
std::ofstream output(file_path, std::ios::binary);
EXPECT_TRUE(output.is_open());
@@ -271,6 +281,42 @@ TEST(IcebergDeleteFileReaderHelperTest,
DeletionVectorCacheKeyEscapesPathBoundar
build_iceberg_deletion_vector_cache_key(second_data_file_path,
second_delete_file));
}
+TEST(IcebergDeleteFileReaderHelperTest, ValidateDeletionVectorDescriptor) {
+ size_t bytes_read = 0;
+
+ TIcebergDeleteFileDesc missing_path;
+ missing_path.__set_content_offset(0);
+ missing_path.__set_content_size_in_bytes(12);
+ EXPECT_FALSE(validate_iceberg_deletion_vector_descriptor(missing_path,
bytes_read).ok());
+
+ EXPECT_FALSE(validate_iceberg_deletion_vector_descriptor(
+ make_iceberg_deletion_vector("dv.puffin", -1, 12),
bytes_read)
+ .ok());
+ EXPECT_FALSE(validate_iceberg_deletion_vector_descriptor(
+ make_iceberg_deletion_vector("dv.puffin", 0, 11),
bytes_read)
+ .ok());
+ EXPECT_TRUE(
+ validate_iceberg_deletion_vector_descriptor(
+ make_iceberg_deletion_vector("dv.puffin", 0,
MAX_ICEBERG_DELETION_VECTOR_BYTES),
+ bytes_read)
+ .ok());
+ EXPECT_EQ(static_cast<size_t>(MAX_ICEBERG_DELETION_VECTOR_BYTES),
bytes_read);
+ const auto unsupported_status =
validate_iceberg_deletion_vector_descriptor(
+ make_iceberg_deletion_vector("dv.puffin", 0,
MAX_ICEBERG_DELETION_VECTOR_BYTES + 1),
+ bytes_read);
+ EXPECT_TRUE(unsupported_status.is<ErrorCode::NOT_IMPLEMENTED_ERROR>()) <<
unsupported_status;
+ EXPECT_FALSE(validate_iceberg_deletion_vector_descriptor(
+ make_iceberg_deletion_vector("dv.puffin",
+
std::numeric_limits<int64_t>::max() - 10, 12),
+ bytes_read)
+ .ok());
+
+ EXPECT_TRUE(validate_iceberg_deletion_vector_descriptor(
+ make_iceberg_deletion_vector("dv.puffin", 7, 12),
bytes_read)
+ .ok());
+ EXPECT_EQ(bytes_read, 12);
+}
+
TEST(IcebergDeleteFileReaderHelperTest, ReadDeletionVectorReportsMissingFile) {
const auto test_dir =
std::filesystem::temp_directory_path() /
"doris_iceberg_deletion_vector_missing_test";
@@ -296,9 +342,9 @@ TEST(IcebergDeleteFileReaderHelperTest,
ReadDeletionVectorReportsMissingFile) {
std::filesystem::remove_all(test_dir);
}
-TEST(IcebergDeleteFileReaderHelperTest, ReadDeletionVectorReportsShortRead) {
+TEST(IcebergDeleteFileReaderHelperTest,
ReadDeletionVectorRejectsRangePastFile) {
const auto test_dir = std::filesystem::temp_directory_path() /
- "doris_iceberg_deletion_vector_short_read_test";
+ "doris_iceberg_deletion_vector_range_past_file_test";
std::filesystem::remove_all(test_dir);
std::filesystem::create_directories(test_dir);
@@ -316,11 +362,49 @@ TEST(IcebergDeleteFileReaderHelperTest,
ReadDeletionVectorReportsShortRead) {
auto status = read_iceberg_deletion_vector(
make_iceberg_deletion_vector(dv_path, 0, dv_size + 1), options,
&rows_to_delete);
- EXPECT_FALSE(status.ok());
+ EXPECT_TRUE(status.is<ErrorCode::DATA_QUALITY_ERROR>()) << status;
+ EXPECT_NE(status.to_string().find("range exceeds file size"),
std::string::npos);
+ EXPECT_NE(status.to_string().find(dv_path), std::string::npos);
EXPECT_EQ(rows_to_delete.cardinality(), 0);
std::filesystem::remove_all(test_dir);
}
+TEST(IcebergDeleteFileReaderHelperTest,
DeletionVectorReaderValidatesOpenedFileRange) {
+ const auto test_dir =
+ std::filesystem::temp_directory_path() /
"doris_deletion_vector_opened_file_range_test";
+ std::filesystem::remove_all(test_dir);
+ std::filesystem::create_directories(test_dir);
+
+ const auto dv_path = (test_dir / "delete-vector.bin").string();
+ const auto dv_size = write_iceberg_deletion_vector_file(dv_path, {1, 3});
+
+ RuntimeProfile profile("test_profile");
+ RuntimeState state {TQueryOptions(), TQueryGlobals()};
+ auto scan_params = make_local_parquet_scan_params();
+ IcebergDeleteFileIOContext io_context(&state);
+
+ {
+ TFileRangeDesc exact_range = build_iceberg_delete_file_range(dv_path);
+ exact_range.start_offset = 4;
+ exact_range.size = dv_size - exact_range.start_offset;
+ DeletionVectorReader exact_reader(&state, &profile, scan_params,
exact_range,
+ &io_context.io_ctx);
+ const auto exact_status = exact_reader.open();
+ EXPECT_TRUE(exact_status.ok()) << exact_status;
+
+ TFileRangeDesc oversized_range = exact_range;
+ oversized_range.size = MAX_ICEBERG_DELETION_VECTOR_BYTES;
+ DeletionVectorReader oversized_reader(&state, &profile, scan_params,
oversized_range,
+ &io_context.io_ctx);
+ const auto oversized_status = oversized_reader.open();
+ EXPECT_TRUE(oversized_status.is<ErrorCode::DATA_QUALITY_ERROR>()) <<
oversized_status;
+ EXPECT_NE(oversized_status.to_string().find("range exceeds file
size"), std::string::npos);
+ EXPECT_NE(oversized_status.to_string().find(dv_path),
std::string::npos);
+ }
+
+ std::filesystem::remove_all(test_dir);
+}
+
TEST(IcebergDeleteFileReaderHelperTest,
ReadDeletionVectorStopsWhenIoContextStops) {
const auto test_dir =
std::filesystem::temp_directory_path() /
"doris_iceberg_deletion_vector_stop_test";
@@ -362,6 +446,19 @@ TEST(IcebergDeleteFileReaderHelperTest,
DecodeDeletionVectorRejectsCorruptPayloa
EXPECT_EQ(rows_to_delete.cardinality(), 0);
}
+TEST(IcebergDeleteFileReaderHelperTest,
DecodeDeletionVectorRejectsCrcMismatch) {
+ auto corrupted = build_iceberg_deletion_vector_blob({1, 3});
+ corrupted.back() ^= 1;
+
+ roaring::Roaring64Map rows_to_delete;
+ auto status = decode_iceberg_deletion_vector_buffer(corrupted.data(),
corrupted.size(),
+ &rows_to_delete);
+
+ EXPECT_TRUE(status.is<ErrorCode::DATA_QUALITY_ERROR>());
+ EXPECT_NE(status.to_string().find("CRC32 mismatch"), std::string::npos);
+ EXPECT_EQ(rows_to_delete.cardinality(), 0);
+}
+
TEST(IcebergDeleteFileReaderHelperTest,
ReadDeletionVectorReadsMillionDeletePositions) {
const auto test_dir =
std::filesystem::temp_directory_path() /
"doris_iceberg_deletion_vector_large_test";
diff --git a/be/test/format/table/paimon_cpp_reader_test.cpp
b/be/test/format/table/paimon_cpp_reader_test.cpp
index e323f4f1af4..f2602c34fd0 100644
--- a/be/test/format/table/paimon_cpp_reader_test.cpp
+++ b/be/test/format/table/paimon_cpp_reader_test.cpp
@@ -21,12 +21,14 @@
#include <gtest/gtest.h>
#include <cstring>
+#include <limits>
#include <string>
#include <vector>
#include "core/block/block.h"
#include "exec/common/endian.h"
#include "format/format_common.h"
+#include "format/table/deletion_vector_reader.h"
#include "format/table/paimon_reader.h"
#include "io/fs/file_meta_cache.h"
#include "io/io_common.h"
@@ -172,9 +174,19 @@ TEST(PaimonDeletionVectorTest, RejectShortBuffer) {
decode_paimon_deletion_vector_buffer(buffer.data(), buffer.size(),
&deletion_vector);
ASSERT_FALSE(status.ok());
+ EXPECT_TRUE(status.is<ErrorCode::DATA_QUALITY_ERROR>());
EXPECT_NE(status.to_string().find("file size too small"),
std::string::npos);
}
+TEST(PaimonDeletionVectorTest, RejectNullBuffer) {
+ DeletionVector deletion_vector;
+ const auto status = decode_paimon_deletion_vector_buffer(nullptr, 8,
&deletion_vector);
+
+ ASSERT_FALSE(status.ok());
+ EXPECT_TRUE(status.is<ErrorCode::DATA_QUALITY_ERROR>());
+ EXPECT_NE(status.to_string().find("blob is null"), std::string::npos);
+}
+
TEST(PaimonDeletionVectorTest, RejectLengthMismatch) {
// Scenario: the big-endian length prefix protects against using a
truncated or over-read DV
// slice from a shared deletion-vector file.
@@ -185,7 +197,8 @@ TEST(PaimonDeletionVectorTest, RejectLengthMismatch) {
decode_paimon_deletion_vector_buffer(buffer.data(), buffer.size(),
&deletion_vector);
ASSERT_FALSE(status.ok());
- EXPECT_NE(status.to_string().find("length not match"), std::string::npos);
+ EXPECT_TRUE(status.is<ErrorCode::DATA_QUALITY_ERROR>());
+ EXPECT_NE(status.to_string().find("length mismatch"), std::string::npos);
}
TEST(PaimonDeletionVectorTest, RejectMagicMismatch) {
@@ -198,7 +211,8 @@ TEST(PaimonDeletionVectorTest, RejectMagicMismatch) {
decode_paimon_deletion_vector_buffer(buffer.data(), buffer.size(),
&deletion_vector);
ASSERT_FALSE(status.ok());
- EXPECT_NE(status.to_string().find("invalid magic number"),
std::string::npos);
+ EXPECT_TRUE(status.is<ErrorCode::DATA_QUALITY_ERROR>());
+ EXPECT_NE(status.to_string().find("magic number mismatch"),
std::string::npos);
}
TEST(PaimonDeletionVectorTest, RejectCorruptRoaringBitmap) {
@@ -234,6 +248,40 @@ TEST(PaimonDeletionVectorTest,
CacheKeyIncludesOffsetAndLength) {
EXPECT_NE(first_key,
build_paimon_deletion_vector_cache_key(different_length));
}
+TEST(PaimonDeletionVectorTest, ValidateDescriptorRejectsInvalidRange) {
+ size_t bytes_read = 0;
+
+ TPaimonDeletionFileDesc missing_path;
+ missing_path.__set_offset(0);
+ missing_path.__set_length(4);
+ EXPECT_FALSE(validate_paimon_deletion_vector_descriptor(missing_path,
bytes_read).ok());
+
+ TPaimonDeletionFileDesc deletion_file;
+ deletion_file.__set_path("dv.bin");
+ deletion_file.__set_offset(-1);
+ deletion_file.__set_length(4);
+ EXPECT_FALSE(validate_paimon_deletion_vector_descriptor(deletion_file,
bytes_read).ok());
+
+ deletion_file.__set_offset(0);
+ deletion_file.__set_length(-1);
+ EXPECT_FALSE(validate_paimon_deletion_vector_descriptor(deletion_file,
bytes_read).ok());
+
+ deletion_file.__set_length(std::numeric_limits<int64_t>::max());
+ EXPECT_FALSE(validate_paimon_deletion_vector_descriptor(deletion_file,
bytes_read).ok());
+
+ deletion_file.__set_length(MAX_PAIMON_DELETION_VECTOR_BYTES - 4);
+ EXPECT_TRUE(validate_paimon_deletion_vector_descriptor(deletion_file,
bytes_read).ok());
+ EXPECT_EQ(static_cast<size_t>(MAX_PAIMON_DELETION_VECTOR_BYTES),
bytes_read);
+
+ deletion_file.__set_length(MAX_PAIMON_DELETION_VECTOR_BYTES - 3);
+ EXPECT_FALSE(validate_paimon_deletion_vector_descriptor(deletion_file,
bytes_read).ok());
+
+ deletion_file.__set_offset(3);
+ deletion_file.__set_length(4);
+ EXPECT_TRUE(validate_paimon_deletion_vector_descriptor(deletion_file,
bytes_read).ok());
+ EXPECT_EQ(bytes_read, 8);
+}
+
TEST(PaimonDeletionVectorTest, DecodedCacheReportsHitSeparatelyFromFileCache) {
// The decoded cache lookup result is reported by ShardedKVCache itself.
The creator represents
// the lower File Cache/read/decode path and must only run for the miss.
diff --git a/be/test/format_v2/table/iceberg_reader_test.cpp
b/be/test/format_v2/table/iceberg_reader_test.cpp
index cbcf6d17256..71dc1ea1db7 100644
--- a/be/test/format_v2/table/iceberg_reader_test.cpp
+++ b/be/test/format_v2/table/iceberg_reader_test.cpp
@@ -73,6 +73,7 @@
#include "runtime/runtime_state.h"
#include "storage/segment/condition_cache.h"
#include "util/debug_points.h"
+#include "util/hash_util.hpp"
namespace doris::format {
namespace {
@@ -713,7 +714,8 @@ int64_t write_iceberg_deletion_vector_file(const
std::string& file_path,
BigEndian::Store32(blob.data(), total_length);
constexpr char DV_MAGIC[] = {'\xD1', '\xD3', '\x39', '\x64'};
memcpy(blob.data() + 4, DV_MAGIC, 4);
- BigEndian::Store32(blob.data() + 8 + bitmap_size, 0);
+ const uint32_t crc = HashUtil::zlib_crc_hash(blob.data() + 4,
total_length, 0);
+ BigEndian::Store32(blob.data() + 8 + bitmap_size, crc);
std::ofstream output(file_path, std::ios::binary);
EXPECT_TRUE(output.is_open());
@@ -1607,8 +1609,25 @@ TEST(IcebergV2ReaderTest,
IcebergDeletionVectorRejectsMissingRange) {
auto status = reader.parse_deletion_vector_file(table_format_desc, &desc,
&has_delete_file);
EXPECT_FALSE(status.ok());
- EXPECT_TRUE(status.is<ErrorCode::INTERNAL_ERROR>());
- EXPECT_NE(status.to_string().find("missing content offset or length"),
std::string::npos);
+ EXPECT_TRUE(status.is<ErrorCode::DATA_QUALITY_ERROR>());
+ EXPECT_NE(status.to_string().find("descriptor misses"), std::string::npos);
+ EXPECT_FALSE(has_delete_file);
+}
+
+TEST(IcebergV2ReaderTest, IcebergDeletionVectorRejectsInvalidRange) {
+ TTableFormatFileDesc table_format_desc;
+ TIcebergFileDesc iceberg_desc;
+ iceberg_desc.__set_format_version(2);
+ iceberg_desc.__set_delete_files({make_iceberg_deletion_vector("dv.bin",
-1, 12)});
+ table_format_desc.__set_iceberg_params(iceberg_desc);
+
+ IcebergTableReaderDeleteFileTestHelper reader;
+ DeleteFileDesc desc;
+ bool has_delete_file = false;
+ auto status = reader.parse_deletion_vector_file(table_format_desc, &desc,
&has_delete_file);
+
+ EXPECT_TRUE(status.is<ErrorCode::DATA_QUALITY_ERROR>());
+ EXPECT_NE(status.to_string().find("offset must be non-negative"),
std::string::npos);
EXPECT_FALSE(has_delete_file);
}
diff --git a/be/test/format_v2/table/paimon_reader_test.cpp
b/be/test/format_v2/table/paimon_reader_test.cpp
index ae1691a774f..fbb3c79d561 100644
--- a/be/test/format_v2/table/paimon_reader_test.cpp
+++ b/be/test/format_v2/table/paimon_reader_test.cpp
@@ -461,6 +461,20 @@ TEST(PaimonReaderTest,
DeletionVectorCacheKeyIncludesOffsetAndLength) {
EXPECT_NE(first_desc.key, different_length_desc.key);
}
+TEST(PaimonReaderTest, DeletionVectorRejectsInvalidRange) {
+ auto table_format_params = make_paimon_table_format_desc("dv.bin", -1, 4);
+
+ paimon::PaimonReader reader;
+ DeleteFileDesc desc;
+ bool has_delete_file = false;
+ auto status =
+ reader.TEST_parse_deletion_vector_file(table_format_params, &desc,
&has_delete_file);
+
+ EXPECT_TRUE(status.is<ErrorCode::DATA_QUALITY_ERROR>());
+ EXPECT_NE(status.to_string().find("offset must be non-negative"),
std::string::npos);
+ EXPECT_FALSE(has_delete_file);
+}
+
TEST(PaimonReaderTest, DecodeDeletionVectorBufferUsesSharedFormatHelper) {
// Scenario: format_v2 TableReader reads a raw Paimon BitmapDeletionVector
range and delegates
// the binary parsing to the same helper used by the format reader path.
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergDeleteFileFilter.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergDeleteFileFilter.java
index 32de4ebfdd9..083409adda1 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergDeleteFileFilter.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergDeleteFileFilter.java
@@ -54,16 +54,46 @@ public class IcebergDeleteFileFilter {
String deleteFilePath = deleteFile.path().toString();
if (deleteFile.format() == FileFormat.PUFFIN) {
+ long fileSize = deleteFile.fileSizeInBytes();
+ Long contentOffset = deleteFile.contentOffset();
+ Long contentLength = deleteFile.contentSizeInBytes();
+ validateDeletionVectorMetadata(deleteFilePath, fileSize,
contentOffset, contentLength);
// The content_offset and content_size_in_bytes fields are used to
reference
// a specific blob for direct access to a deletion vector.
return new DeletionVector(deleteFilePath,
positionLowerBound.orElse(-1L), positionUpperBound.orElse(-1L),
- deleteFile.fileSizeInBytes(), deleteFile.contentOffset(),
deleteFile.contentSizeInBytes());
+ fileSize, contentOffset, contentLength);
} else {
return new PositionDelete(deleteFilePath,
positionLowerBound.orElse(-1L), positionUpperBound.orElse(-1L),
deleteFile.fileSizeInBytes(), deleteFile.format());
}
}
+ static void validateDeletionVectorMetadata(
+ String deleteFilePath, long fileSize, Long contentOffset, Long
contentLength) {
+ if (contentOffset == null || contentLength == null) {
+ throw new IllegalArgumentException(String.format(
+ "Iceberg deletion vector metadata misses content offset or
length: %s", deleteFilePath));
+ }
+ if (fileSize < 0 || contentOffset < 0 || contentLength < 0) {
+ throw new IllegalArgumentException(String.format(
+ "Iceberg deletion vector metadata must be non-negative,
file: %s, file size: %d, "
+ + "content offset: %d, content length: %d",
+ deleteFilePath, fileSize, contentOffset, contentLength));
+ }
+ if (contentOffset > Long.MAX_VALUE - contentLength) {
+ throw new IllegalArgumentException(String.format(
+ "Iceberg deletion vector metadata range overflows, file:
%s, content offset: %d, "
+ + "content length: %d",
+ deleteFilePath, contentOffset, contentLength));
+ }
+ if (contentOffset + contentLength > fileSize) {
+ throw new IllegalArgumentException(String.format(
+ "Iceberg deletion vector metadata range exceeds file size,
file: %s, file size: %d, "
+ + "content offset: %d, content length: %d",
+ deleteFilePath, fileSize, contentOffset, contentLength));
+ }
+ }
+
public static EqualityDelete createEqualityDelete(String deleteFilePath,
List<Integer> fieldIds,
long fileSize, FileFormat fileformat) {
// todo:
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
index 74de233b65b..eee7bc40a92 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
@@ -1076,10 +1076,14 @@ public class IcebergScanNode extends FileQueryScanNode {
split.setPositionDeleteFileFormat(getNativePositionDeleteFileFormat(deleteFile.format()));
split.setPositionDeleteOriginalPath(originalPath);
if (deleteFile.format() == FileFormat.PUFFIN) {
+ Long contentOffset = deleteFile.contentOffset();
+ Long contentLength = deleteFile.contentSizeInBytes();
+ IcebergDeleteFileFilter.validateDeletionVectorMetadata(
+ originalPath, deleteFile.fileSizeInBytes(), contentOffset,
contentLength);
split.setPositionDeleteContent(IcebergDeleteFileFilter.DeletionVector.type());
split.setPositionDeleteReferencedDataFilePath(deleteFile.referencedDataFile());
- split.setPositionDeleteContentOffset(deleteFile.contentOffset());
-
split.setPositionDeleteContentSizeInBytes(deleteFile.contentSizeInBytes());
+ split.setPositionDeleteContentOffset(contentOffset);
+ split.setPositionDeleteContentSizeInBytes(contentLength);
} else {
split.setPositionDeleteContent(IcebergDeleteFileFilter.PositionDelete.type());
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergDeleteFileFilterTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergDeleteFileFilterTest.java
new file mode 100644
index 00000000000..6a714107a6e
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergDeleteFileFilterTest.java
@@ -0,0 +1,48 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.datasource.iceberg.source;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+public class IcebergDeleteFileFilterTest {
+ @Test
+ public void testValidateDeletionVectorMetadataAcceptsLongRange() {
+ Assertions.assertDoesNotThrow(() ->
IcebergDeleteFileFilter.validateDeletionVectorMetadata(
+ "puffin.dv", 1L << 40, (1L << 32) + 17, (1L << 30) + 19));
+ }
+
+ @Test
+ public void testValidateDeletionVectorMetadataRejectsInvalidRange() {
+ Assertions.assertThrows(IllegalArgumentException.class,
+ () ->
IcebergDeleteFileFilter.validateDeletionVectorMetadata("puffin.dv", 100, null,
1L));
+ Assertions.assertThrows(IllegalArgumentException.class,
+ () ->
IcebergDeleteFileFilter.validateDeletionVectorMetadata("puffin.dv", 100, 1L,
null));
+ Assertions.assertThrows(IllegalArgumentException.class,
+ () ->
IcebergDeleteFileFilter.validateDeletionVectorMetadata("puffin.dv", -1, 1L,
1L));
+ Assertions.assertThrows(IllegalArgumentException.class,
+ () ->
IcebergDeleteFileFilter.validateDeletionVectorMetadata("puffin.dv", 100, -1L,
1L));
+ Assertions.assertThrows(IllegalArgumentException.class,
+ () ->
IcebergDeleteFileFilter.validateDeletionVectorMetadata("puffin.dv", 100, 1L,
-1L));
+ Assertions.assertThrows(IllegalArgumentException.class,
+ () -> IcebergDeleteFileFilter.validateDeletionVectorMetadata(
+ "puffin.dv", Long.MAX_VALUE, Long.MAX_VALUE, 1L));
+ Assertions.assertThrows(IllegalArgumentException.class,
+ () ->
IcebergDeleteFileFilter.validateDeletionVectorMetadata("puffin.dv", 100, 90L,
11L));
+ }
+}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
index 16aaf739a91..db0f52cd871 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
@@ -36,10 +36,12 @@ import org.apache.doris.thrift.TIcebergDeleteFileDesc;
import org.apache.doris.thrift.TPushAggOp;
import org.apache.iceberg.DataFile;
+import org.apache.iceberg.DeleteFile;
import org.apache.iceberg.FileFormat;
import org.apache.iceberg.FileScanTask;
import org.apache.iceberg.PartitionData;
import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.PositionDeletesScanTask;
import org.apache.iceberg.Schema;
import org.apache.iceberg.Snapshot;
import org.apache.iceberg.Table;
@@ -343,6 +345,34 @@ public class IcebergScanNodeTest {
Assert.assertEquals(TFileFormatType.FORMAT_ORC,
rangeDesc.getFormatType());
}
+ @Test
+ public void testPositionDeleteSystemTableValidatesDeletionVectorMetadata()
throws Exception {
+ DeleteFile deleteFile = Mockito.mock(DeleteFile.class);
+
Mockito.when(deleteFile.path()).thenReturn("file:///tmp/delete-shared.puffin");
+ Mockito.when(deleteFile.format()).thenReturn(FileFormat.PUFFIN);
+ Mockito.when(deleteFile.fileSizeInBytes()).thenReturn(100L);
+ Mockito.when(deleteFile.contentOffset()).thenReturn(null);
+ Mockito.when(deleteFile.contentSizeInBytes()).thenReturn(10L);
+
+ PositionDeletesScanTask task =
Mockito.mock(PositionDeletesScanTask.class);
+ Mockito.when(task.file()).thenReturn(deleteFile);
+ Mockito.when(task.start()).thenReturn(0L);
+ Mockito.when(task.length()).thenReturn(100L);
+
+ TestIcebergScanNode node = new TestIcebergScanNode(new
SessionVariable());
+ Method method = IcebergScanNode.class.getDeclaredMethod(
+ "createIcebergPositionDeleteSysSplit",
PositionDeletesScanTask.class);
+ method.setAccessible(true);
+
+ try {
+ method.invoke(node, task);
+ Assert.fail("position_deletes planning should reject invalid
deletion vector metadata");
+ } catch (InvocationTargetException e) {
+ Assert.assertTrue(e.getCause() instanceof
IllegalArgumentException);
+
Assert.assertTrue(e.getCause().getMessage().contains("delete-shared.puffin"));
+ }
+ }
+
@Test
public void testSetIcebergParamsPropagatesPositionDeleteFileFormat()
throws Exception {
SessionVariable sv = new SessionVariable();
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]