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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -484,6 +484,45 @@ TEST_F(LookupMergeTreeCompactRewriterTest, TestFirstRowRewrite) {
CheckResult(compact_file_name, table_schema, "orc", expected_array);
}

TEST_F(LookupMergeTreeCompactRewriterTest, TestFirstRowLooksUpExistingKeys) {
std::map<std::string, std::string> options = {{Options::MERGE_ENGINE, "first-row"},
{Options::FILE_FORMAT, "orc"}};
ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options));
ASSERT_OK_AND_ASSIGN(auto table_path, CreateTable(options));
auto schema_manager = std::make_shared<SchemaManager>(fs_, table_path);
ASSERT_OK_AND_ASSIGN(auto table_schema, schema_manager->ReadSchema(0));

ASSERT_OK_AND_ASSIGN(auto level0_file,
NewFiles(/*level=*/0, /*last_sequence_number=*/0, table_path, core_options,
"[[1, 111], [2, 22]]"));
ASSERT_OK_AND_ASSIGN(auto high_level_file, NewFiles(/*level=*/2, /*last_sequence_number=*/-1,
table_path, core_options, "[[1, 11]]"));
auto processor_factory = std::make_shared<PersistEmptyProcessor::Factory>();
ASSERT_OK_AND_ASSIGN(auto lookup_levels,
CreateLookupLevels<bool>(table_path, table_schema, processor_factory,
std::vector<std::shared_ptr<DataFileMeta>>{
level0_file, high_level_file}));
ASSERT_OK_AND_ASSIGN(auto rewriter,
CreateCompactRewriterForFirstRow(table_path, table_schema, core_options,
std::move(lookup_levels)));
ASSERT_OK_AND_ASSIGN(
auto runs, GenerateSortedRuns(std::vector<std::shared_ptr<DataFileMeta>>{level0_file}));
ASSERT_OK_AND_ASSIGN(auto compact_result, rewriter->Rewrite(
/*output_level=*/1, /*drop_delete=*/true, runs));

ASSERT_EQ(1, compact_result.After().size());
ASSERT_EQ(1, compact_result.After()[0]->row_count);

auto type_with_special_fields =
arrow::struct_(SpecialFields::CompleteSequenceAndValueKindField(arrow_schema_)->fields());
std::shared_ptr<arrow::ChunkedArray> expected;
ASSERT_TRUE(arrow::ipc::internal::json::ChunkedArrayFromJSON(type_with_special_fields,
{"[[2, 0, 2, 22]]"}, &expected)
.ok());
CheckResult(table_path + "/bucket-0/" + compact_result.After()[0]->file_name, table_schema,
"orc", expected);
}

TEST_F(LookupMergeTreeCompactRewriterTest, TestFirstRowUpgrade) {
std::map<std::string, std::string> options = {{Options::MERGE_ENGINE, "first-row"},
{Options::FILE_FORMAT, "orc"}};
Expand Down
14 changes: 2 additions & 12 deletions src/paimon/core/mergetree/lookup_levels.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -162,8 +162,8 @@ LookupLevels<T>::LookupLevels(
lookup_store_factory_(lookup_store_factory),
lookup_file_cache_(lookup_file_cache),
remote_lookup_file_manager_(remote_lookup_file_manager) {
if constexpr (std::is_same_v<T, FilePosition>) {
// if T is FilePosition, only read key fields to create sst file is enough
if constexpr (std::is_same_v<T, FilePosition> || std::is_same_v<T, bool>) {
// FilePosition and first-row lookup do not persist values, so reading key fields is enough.
value_schema_ = key_schema_;
} else {
value_schema_ = DataField::ConvertDataFieldsToArrowSchema(table_schema->Fields());
Expand Down Expand Up @@ -332,16 +332,6 @@ std::optional<std::string> LookupLevels<T>::TryToDownloadRemoteSst(
template <typename T>
Status LookupLevels<T>::CreateSstFileFromDataFile(const std::shared_ptr<DataFileMeta>& file,
const std::string& kv_file_path) {
if constexpr (std::is_same_v<T, bool>) {
// Short-circuit logic: if T is bool, just write empty lookup file.
PAIMON_ASSIGN_OR_RAISE(
std::shared_ptr<BloomFilter> bloom_filter,
LookupStoreFactory::BfGenerator(file->row_count, options_, pool_.get()));
PAIMON_ASSIGN_OR_RAISE(
std::unique_ptr<LookupStoreWriter> kv_writer,
lookup_store_factory_->CreateWriter(fs_, kv_file_path, bloom_filter, pool_));
return kv_writer->Close();
}
// Prepare reader to iterate KeyValue
PAIMON_ASSIGN_OR_RAISE(
std::vector<std::unique_ptr<FileBatchReader>> raw_readers,
Expand Down
Loading