Skip to content

Commit 65bc021

Browse files
committed
[fix](iceberg) Add missing Iceberg field IDs for position delete files
1 parent f852097 commit 65bc021

3 files changed

Lines changed: 72 additions & 4 deletions

File tree

be/src/exec/sink/writer/iceberg/viceberg_delete_file_writer.cpp

Lines changed: 21 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -19,13 +19,29 @@
1919

2020
#include <fmt/format.h>
2121

22+
#include "format/table/iceberg/schema.h"
23+
#include "format/table/iceberg/types.h"
2224
#include "format/transformer/vorc_transformer.h"
2325
#include "format/transformer/vparquet_transformer.h"
2426
#include "io/file_factory.h"
2527
#include "runtime/runtime_state.h"
2628

2729
namespace doris {
2830

31+
// Iceberg reserved field IDs for position delete files.
32+
constexpr int POSITION_DELETE_FILE_PATH_ID = 2147483546;
33+
constexpr int POSITION_DELETE_POS_ID = 2147483545;
34+
35+
std::unique_ptr<iceberg::Schema> build_position_delete_schema() {
36+
std::vector<iceberg::NestedField> fields;
37+
fields.reserve(2);
38+
fields.emplace_back(false, POSITION_DELETE_FILE_PATH_ID, "file_path",
39+
std::make_unique<iceberg::StringType>(), std::nullopt);
40+
fields.emplace_back(false, POSITION_DELETE_POS_ID, "pos", std::make_unique<iceberg::LongType>(),
41+
std::nullopt);
42+
return std::make_unique<iceberg::Schema>(std::move(fields));
43+
}
44+
2945
VIcebergDeleteFileWriter::VIcebergDeleteFileWriter(TFileContent::type delete_type,
3046
const std::string& output_path,
3147
TFileFormatType::type file_format,
@@ -46,6 +62,7 @@ Status VIcebergDeleteFileWriter::open(RuntimeState* state, RuntimeProfile* profi
4662
if (_delete_type != TFileContent::POSITION_DELETES) {
4763
return Status::NotSupported("Iceberg delete file writer only supports position deletes");
4864
}
65+
_position_delete_schema = build_position_delete_schema();
4966

5067
_state = state;
5168

@@ -83,15 +100,15 @@ Status VIcebergDeleteFileWriter::open(RuntimeState* state, RuntimeProfile* profi
83100

84101
ParquetFileOptions parquet_options = {parquet_compression_type,
85102
TParquetVersion::PARQUET_1_0, false, false};
86-
_file_format_transformer.reset(new VParquetTransformer(state, _file_writer.get(),
87-
output_exprs, column_names, false,
88-
parquet_options, nullptr, nullptr));
103+
_file_format_transformer.reset(new VParquetTransformer(
104+
state, _file_writer.get(), output_exprs, column_names, false, parquet_options,
105+
nullptr, _position_delete_schema.get()));
89106
return _file_format_transformer->open();
90107
}
91108
case TFileFormatType::FORMAT_ORC: {
92109
_file_format_transformer.reset(new VOrcTransformer(state, _file_writer.get(), output_exprs,
93110
"", column_names, false, _compress_type,
94-
nullptr));
111+
_position_delete_schema.get()));
95112
return _file_format_transformer->open();
96113
}
97114
default:

be/src/exec/sink/writer/iceberg/viceberg_delete_file_writer.h

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,9 @@ namespace doris {
3232
class RuntimeState;
3333
class RuntimeProfile;
3434
class ObjectPool;
35+
namespace iceberg {
36+
class Schema;
37+
}
3538

3639
namespace io {
3740
class FileSystem;
@@ -103,6 +106,7 @@ class VIcebergDeleteFileWriter {
103106
RuntimeState* _state = nullptr;
104107
std::shared_ptr<io::FileSystem> _fs;
105108
io::FileWriterPtr _file_writer;
109+
std::unique_ptr<iceberg::Schema> _position_delete_schema;
106110
std::unique_ptr<VFileFormatTransformer> _file_format_transformer;
107111

108112
int32_t _partition_spec_id = 0;

be/test/exec/sink/viceberg_delete_sink_test.cpp

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,8 @@
1818
#include "exec/sink/viceberg_delete_sink.h"
1919

2020
#include <gtest/gtest.h>
21+
#include <parquet/api/reader.h>
22+
#include <parquet/schema.h>
2123
#include <rapidjson/document.h>
2224

2325
#include <filesystem>
@@ -35,7 +37,9 @@
3537
#include "exec/common/endian.h"
3638
#include "gen_cpp/DataSinks_types.h"
3739
#include "gen_cpp/Types_types.h"
40+
#include "runtime/runtime_profile.h"
3841
#include "runtime/runtime_state.h"
42+
#include "testutil/mock/mock_runtime_state.h"
3943
#include "util/uid_util.h"
4044

4145
namespace doris {
@@ -480,6 +484,49 @@ TEST_F(VIcebergDeleteSinkTest, TestGenerateDeleteFilePath) {
480484
ASSERT_NE(std::string::npos, delete_file_path.find("delete_pos_"));
481485
}
482486

487+
TEST_F(VIcebergDeleteSinkTest, TestWritePositionDeleteParquetFieldIds) {
488+
std::filesystem::path temp_dir = std::filesystem::temp_directory_path() /
489+
("iceberg_position_delete_test_" + generate_uuid_string());
490+
ASSERT_TRUE(std::filesystem::create_directories(temp_dir));
491+
492+
TDataSink t_data_sink = build_local_delete_sink(temp_dir.string(), 2);
493+
VExprContextSPtrs output_exprs;
494+
auto sink = std::make_shared<VIcebergDeleteSink>(t_data_sink, output_exprs, nullptr, nullptr);
495+
ObjectPool pool;
496+
ASSERT_TRUE(sink->init_properties(&pool).ok());
497+
498+
MockRuntimeState state;
499+
RuntimeProfile profile("iceberg_delete_sink");
500+
ASSERT_TRUE(sink->open(&state, &profile).ok());
501+
502+
std::map<std::string, IcebergFileDeletion> file_deletions;
503+
auto [file_it, inserted] =
504+
file_deletions.emplace("file1.parquet", IcebergFileDeletion(1, "[\"p=1\"]"));
505+
ASSERT_TRUE(inserted);
506+
file_it->second.rows_to_delete.add((uint32_t)10);
507+
file_it->second.rows_to_delete.add((uint32_t)20);
508+
509+
ASSERT_TRUE(sink->_write_position_delete_files(file_deletions).ok());
510+
ASSERT_EQ(1, sink->_commit_data_list.size());
511+
512+
const auto& commit_data = sink->_commit_data_list[0];
513+
std::unique_ptr<::parquet::ParquetFileReader> parquet_reader =
514+
::parquet::ParquetFileReader::OpenFile(commit_data.file_path, false);
515+
std::shared_ptr<::parquet::FileMetaData> file_metadata = parquet_reader->metadata();
516+
const auto& group_node = static_cast<const ::parquet::schema::GroupNode&>(
517+
*file_metadata->schema()->group_node());
518+
519+
ASSERT_EQ(2, group_node.field_count());
520+
auto file_path_field = group_node.field(0);
521+
auto pos_field = group_node.field(1);
522+
EXPECT_EQ("file_path", file_path_field->name());
523+
EXPECT_EQ(2147483546, file_path_field->field_id());
524+
EXPECT_EQ("pos", pos_field->name());
525+
EXPECT_EQ(2147483545, pos_field->field_id());
526+
527+
ASSERT_TRUE(std::filesystem::remove_all(temp_dir) > 0);
528+
}
529+
483530
TEST_F(VIcebergDeleteSinkTest, TestUnsupportedDeleteType) {
484531
// Create a TDataSink for an unsupported delete type
485532
TDataSink t_eq_delete_sink;

0 commit comments

Comments
 (0)