Skip to content

Commit e4c90c4

Browse files
committed
fix comments
1 parent 403876b commit e4c90c4

16 files changed

Lines changed: 466 additions & 82 deletions

src/paimon/CMakeLists.txt

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -882,6 +882,7 @@ if(PAIMON_BUILD_TESTS)
882882
core/utils/file_utils_test.cpp
883883
core/utils/manifest_meta_reader_test.cpp
884884
core/utils/offset_row_test.cpp
885+
core/utils/partition_utils_test.cpp
885886
core/utils/partition_path_utils_test.cpp
886887
core/utils/snapshot_manager_test.cpp
887888
core/utils/tag_manager_test.cpp

src/paimon/common/utils/binary_row_partition_computer.cpp

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,48 @@ Result<BinaryRow> BinaryRowPartitionComputer::ToBinaryRow(
102102
return binary_row;
103103
}
104104

105+
Result<std::map<std::string, std::string>> BinaryRowPartitionComputer::NormalizePartitionSpec(
106+
const std::map<std::string, std::string>& partition) const {
107+
for (const auto& [partition_key, _] : partition) {
108+
if (std::find(partition_keys_.begin(), partition_keys_.end(), partition_key) ==
109+
partition_keys_.end()) {
110+
return Status::Invalid(
111+
fmt::format("field {} does not exist in partition keys", partition_key));
112+
}
113+
}
114+
115+
BinaryRow binary_row(partition_converters_.size());
116+
BinaryRowWriter writer(&binary_row, /*initial_size=*/0, memory_pool_.get());
117+
std::vector<bool> included_fields(partition_converters_.size(), false);
118+
for (size_t field_idx = 0; field_idx < partition_converters_.size(); ++field_idx) {
119+
const PartitionConverter& partition_converter = partition_converters_[field_idx];
120+
auto input_iter = partition.find(partition_converter.partition_key);
121+
if (input_iter == partition.end()) {
122+
writer.SetNullAt(field_idx);
123+
continue;
124+
}
125+
126+
included_fields[field_idx] = true;
127+
if (input_iter->second == default_part_value_) {
128+
writer.SetNullAt(field_idx);
129+
} else {
130+
PAIMON_RETURN_NOT_OK(
131+
partition_converter.converter(input_iter->second, field_idx, &writer));
132+
}
133+
}
134+
writer.Complete();
135+
136+
std::vector<std::pair<std::string, std::string>> normalized_values;
137+
PAIMON_ASSIGN_OR_RAISE(normalized_values, GeneratePartitionVector(binary_row));
138+
std::map<std::string, std::string> normalized_partition;
139+
for (size_t field_idx = 0; field_idx < normalized_values.size(); ++field_idx) {
140+
if (included_fields[field_idx]) {
141+
normalized_partition.insert(normalized_values[field_idx]);
142+
}
143+
}
144+
return normalized_partition;
145+
}
146+
105147
Result<std::vector<std::pair<std::string, std::string>>>
106148
BinaryRowPartitionComputer::GeneratePartitionVector(const BinaryRow& partition) const {
107149
if (static_cast<size_t>(partition.GetFieldCount()) != partition_converters_.size()) {

src/paimon/common/utils/binary_row_partition_computer.h

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,8 @@ class BinaryRowPartitionComputer {
5555
bool legacy_partition_name_enabled, const std::shared_ptr<MemoryPool>& memory_pool);
5656

5757
Result<BinaryRow> ToBinaryRow(const std::map<std::string, std::string>& partition) const;
58+
Result<std::map<std::string, std::string>> NormalizePartitionSpec(
59+
const std::map<std::string, std::string>& partition) const;
5860
Result<std::vector<std::pair<std::string, std::string>>> GeneratePartitionVector(
5961
const BinaryRow& partition) const;
6062
const std::vector<std::string>& GetPartitionKeys() const {

src/paimon/common/utils/binary_row_partition_computer_test.cpp

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -268,6 +268,36 @@ TEST(BinaryRowPartitionComputerTest, TestNullOrWhitespaceOnlyStr) {
268268
ASSERT_EQ(partition_key_values, expected);
269269
}
270270

271+
TEST(BinaryRowPartitionComputerTest, TestNormalizePartialPartitionSpec) {
272+
using PartitionSpec = std::map<std::string, std::string>;
273+
274+
std::shared_ptr<MemoryPool> pool = GetDefaultPool();
275+
std::shared_ptr<arrow::Schema> schema =
276+
arrow::schema({arrow::field("pt", arrow::date32()), arrow::field("region", arrow::utf8())});
277+
std::vector<std::string> partition_keys = {"pt", "region"};
278+
279+
ASSERT_OK_AND_ASSIGN(
280+
std::unique_ptr<BinaryRowPartitionComputer> legacy_computer,
281+
BinaryRowPartitionComputer::Create(partition_keys, schema, "__DEFAULT_PARTITION__",
282+
/*legacy_partition_name_enabled=*/true, pool));
283+
ASSERT_OK_AND_ASSIGN(PartitionSpec legacy_partition,
284+
legacy_computer->NormalizePartitionSpec({{"pt", "2024-01-01"}}));
285+
PartitionSpec expected_legacy_partition = {{"pt", "19723"}};
286+
ASSERT_EQ(expected_legacy_partition, legacy_partition);
287+
288+
ASSERT_OK_AND_ASSIGN(
289+
std::unique_ptr<BinaryRowPartitionComputer> non_legacy_computer,
290+
BinaryRowPartitionComputer::Create(partition_keys, schema, "__DEFAULT_PARTITION__",
291+
/*legacy_partition_name_enabled=*/false, pool));
292+
ASSERT_OK_AND_ASSIGN(PartitionSpec non_legacy_partition,
293+
non_legacy_computer->NormalizePartitionSpec({{"pt", "2024-01-01"}}));
294+
PartitionSpec expected_non_legacy_partition = {{"pt", "2024-01-01"}};
295+
ASSERT_EQ(expected_non_legacy_partition, non_legacy_partition);
296+
297+
ASSERT_NOK_WITH_MSG(non_legacy_computer->NormalizePartitionSpec({{"unknown", "value"}}),
298+
"field unknown does not exist in partition keys");
299+
}
300+
271301
TEST(BinaryRowPartitionComputerTest, TestPartToSimpleString) {
272302
auto pool = GetDefaultPool();
273303
{

src/paimon/core/operation/append_only_file_store_write.cpp

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -106,10 +106,6 @@ Status AppendOnlyFileStoreWrite::RefreshCommittedSnapshot(int64_t snapshot_id) {
106106
return Status::Invalid("refresh committed snapshot requires a real-time writer");
107107
}
108108
PAIMON_ASSIGN_OR_RAISE(Snapshot snapshot, snapshot_manager_->LoadSnapshot(snapshot_id));
109-
if (snapshot.GetCommitKind() == Snapshot::CommitKind::Overwrite()) {
110-
return Status::Invalid(
111-
"real-time committed progress was reset by overwrite; recreate RealtimeContext");
112-
}
113109
PAIMON_ASSIGN_OR_RAISE(
114110
RealtimeOffsetMap committed_offsets,
115111
RealtimeCommitProperties::ReadOffsets(std::optional<Snapshot>(std::move(snapshot)),

src/paimon/core/operation/commit/commit_scanner.cpp

Lines changed: 4 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@
3838
#include "paimon/core/operation/commit/overwrite_changes_provider.h"
3939
#include "paimon/core/operation/file_store_scan.h"
4040
#include "paimon/core/table/bucket_mode.h"
41+
#include "paimon/core/utils/partition_utils.h"
4142
#include "paimon/scan_context.h"
4243

4344
namespace paimon {
@@ -183,14 +184,9 @@ Result<std::vector<IndexManifestEntry>> CommitScanner::ReadAllIndexEntriesFromPa
183184
}
184185

185186
for (const auto& partition_spec : partitions) {
186-
bool matched = true;
187-
for (const auto& [key, value] : partition_spec) {
188-
auto iter = partition.find(key);
189-
if (iter == partition.end() || iter->second != value) {
190-
matched = false;
191-
break;
192-
}
193-
}
187+
PAIMON_ASSIGN_OR_RAISE(bool matched,
188+
PartitionUtils::MatchPartitionSpec(partition, partition_spec,
189+
*partition_computer_));
194190
if (matched) {
195191
return true;
196192
}

src/paimon/core/operation/commit/realtime_commit_properties.cpp

Lines changed: 11 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -31,23 +31,13 @@
3131
#include "paimon/common/utils/rapidjson_util.h"
3232
#include "paimon/common/utils/uuid.h"
3333
#include "paimon/core/utils/branch_manager.h"
34+
#include "paimon/core/utils/partition_utils.h"
3435
#include "paimon/fs/file_system.h"
3536
#include "paimon/macros.h"
3637

3738
namespace paimon {
3839
namespace {
3940

40-
bool MatchPartitionSpec(const std::map<std::string, std::string>& partition,
41-
const std::map<std::string, std::string>& partition_spec) {
42-
for (const auto& [key, value] : partition_spec) {
43-
auto iter = partition.find(key);
44-
if (iter == partition.end() || iter->second != value) {
45-
return false;
46-
}
47-
}
48-
return true;
49-
}
50-
5141
class OffsetEntryJson : public Jsonizable<OffsetEntryJson> {
5242
public:
5343
OffsetEntryJson() = default;
@@ -247,6 +237,7 @@ Result<std::map<std::string, std::string>> RealtimeCommitProperties::Build(
247237
const std::map<RealtimePartitionBucket, OffsetRange>& realtime_ranges,
248238
bool reset_all_realtime_progress,
249239
const std::vector<std::map<std::string, std::string>>& removed_realtime_partitions,
240+
const BinaryRowPartitionComputer& partition_computer,
250241
const std::shared_ptr<FileSystem>& file_system, const std::string& table_root,
251242
const std::string& branch) {
252243
std::map<std::string, std::string> merged_properties = properties;
@@ -269,11 +260,15 @@ Result<std::map<std::string, std::string>> RealtimeCommitProperties::Build(
269260
RealtimeOffsetMap merged_offsets,
270261
ReadOffsets(reset_all_realtime_progress ? std::nullopt : latest_snapshot, file_system));
271262
for (auto iter = merged_offsets.begin(); iter != merged_offsets.end();) {
272-
bool removed =
273-
std::any_of(removed_realtime_partitions.begin(), removed_realtime_partitions.end(),
274-
[&iter](const std::map<std::string, std::string>& partition_spec) {
275-
return MatchPartitionSpec(iter->first.partition, partition_spec);
276-
});
263+
bool removed = false;
264+
for (const auto& partition_spec : removed_realtime_partitions) {
265+
PAIMON_ASSIGN_OR_RAISE(
266+
removed, PartitionUtils::MatchPartitionSpec(iter->first.partition, partition_spec,
267+
partition_computer));
268+
if (removed) {
269+
break;
270+
}
271+
}
277272
if (removed) {
278273
iter = merged_offsets.erase(iter);
279274
} else {

src/paimon/core/operation/commit/realtime_commit_properties.h

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@
3333

3434
namespace paimon {
3535

36+
class BinaryRowPartitionComputer;
3637
class FileSystem;
3738

3839
class RealtimeCommitProperties {
@@ -74,6 +75,7 @@ class RealtimeCommitProperties {
7475
const std::map<RealtimePartitionBucket, OffsetRange>& realtime_ranges,
7576
bool reset_all_realtime_progress,
7677
const std::vector<std::map<std::string, std::string>>& removed_realtime_partitions,
78+
const BinaryRowPartitionComputer& partition_computer,
7779
const std::shared_ptr<FileSystem>& file_system, const std::string& table_root,
7880
const std::string& branch);
7981

src/paimon/core/operation/commit/realtime_commit_properties_test.cpp

Lines changed: 46 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -28,9 +28,12 @@
2828
#include <utility>
2929
#include <vector>
3030

31+
#include "arrow/type.h"
3132
#include "gtest/gtest.h"
33+
#include "paimon/common/utils/binary_row_partition_computer.h"
3234
#include "paimon/fs/file_system.h"
3335
#include "paimon/macros.h"
36+
#include "paimon/memory/memory_pool.h"
3437
#include "paimon/testing/utils/testharness.h"
3538

3639
namespace paimon::test {
@@ -71,6 +74,13 @@ class RealtimeCommitPropertiesTest : public testing::Test {
7174
ASSERT_NE(nullptr, directory_);
7275
file_system_ = directory_->GetFileSystem();
7376
ASSERT_NE(nullptr, file_system_);
77+
std::shared_ptr<arrow::Schema> schema = arrow::schema(
78+
{arrow::field("dt", arrow::utf8()), arrow::field("region", arrow::utf8())});
79+
ASSERT_OK_AND_ASSIGN(partition_computer_,
80+
BinaryRowPartitionComputer::Create(
81+
/*partition_keys=*/{"dt", "region"}, schema,
82+
/*default_part_value=*/"__DEFAULT_PARTITION__",
83+
/*legacy_partition_name_enabled=*/true, GetDefaultPool()));
7484
}
7585

7686
Snapshot MakeSnapshot(
@@ -120,6 +130,7 @@ class RealtimeCommitPropertiesTest : public testing::Test {
120130

121131
std::unique_ptr<UniqueTestDirectory> directory_;
122132
std::shared_ptr<FileSystem> file_system_;
133+
std::unique_ptr<BinaryRowPartitionComputer> partition_computer_;
123134
int32_t next_file_id_ = 0;
124135
};
125136

@@ -312,25 +323,27 @@ TEST_F(RealtimeCommitPropertiesTest, BuildRejectsInvalidProgress) {
312323
ASSERT_NOK_WITH_MSG(
313324
RealtimeCommitProperties::Build(/*properties=*/{}, /*latest_snapshot=*/std::nullopt,
314325
invalid_bucket, /*reset_all_realtime_progress=*/false,
315-
/*removed_realtime_partitions=*/{}, file_system_,
316-
directory_->Str(), "main"),
326+
/*removed_realtime_partitions=*/{}, *partition_computer_,
327+
file_system_, directory_->Str(), "main"),
317328
"bucket -1 is invalid");
318329

319330
std::map<RealtimePartitionBucket, OffsetRange> gap = {
320331
{RealtimePartitionBucket({{"dt", "2"}}, /*bucket=*/0), OffsetRange(3, 5)}};
321-
ASSERT_NOK_WITH_MSG(RealtimeCommitProperties::Build(/*properties=*/{}, latest_snapshot, gap,
322-
/*reset_all_realtime_progress=*/false,
323-
/*removed_realtime_partitions=*/{},
324-
file_system_, directory_->Str(), "main"),
325-
"are not contiguous");
332+
ASSERT_NOK_WITH_MSG(
333+
RealtimeCommitProperties::Build(/*properties=*/{}, latest_snapshot, gap,
334+
/*reset_all_realtime_progress=*/false,
335+
/*removed_realtime_partitions=*/{}, *partition_computer_,
336+
file_system_, directory_->Str(), "main"),
337+
"are not contiguous");
326338

327339
std::map<RealtimePartitionBucket, OffsetRange> overlap = {
328340
{RealtimePartitionBucket({{"dt", "2"}}, /*bucket=*/0), OffsetRange(0, 2)}};
329-
ASSERT_NOK_WITH_MSG(RealtimeCommitProperties::Build(/*properties=*/{}, latest_snapshot, overlap,
330-
/*reset_all_realtime_progress=*/false,
331-
/*removed_realtime_partitions=*/{},
332-
file_system_, directory_->Str(), "main"),
333-
"are not contiguous");
341+
ASSERT_NOK_WITH_MSG(
342+
RealtimeCommitProperties::Build(/*properties=*/{}, latest_snapshot, overlap,
343+
/*reset_all_realtime_progress=*/false,
344+
/*removed_realtime_partitions=*/{}, *partition_computer_,
345+
file_system_, directory_->Str(), "main"),
346+
"are not contiguous");
334347

335348
RealtimeOffsetMap exhausted_offsets = {{bucket0, std::numeric_limits<int64_t>::max()}};
336349
ASSERT_OK_AND_ASSIGN(Snapshot exhausted_snapshot, MakeSnapshotWithOffsets(exhausted_offsets));
@@ -340,8 +353,8 @@ TEST_F(RealtimeCommitPropertiesTest, BuildRejectsInvalidProgress) {
340353
ASSERT_NOK_WITH_MSG(
341354
RealtimeCommitProperties::Build(/*properties=*/{}, exhausted_snapshot, after_max,
342355
/*reset_all_realtime_progress=*/false,
343-
/*removed_realtime_partitions=*/{}, file_system_,
344-
directory_->Str(), "main"),
356+
/*removed_realtime_partitions=*/{}, *partition_computer_,
357+
file_system_, directory_->Str(), "main"),
345358
"offset range is invalid");
346359
}
347360

@@ -355,19 +368,20 @@ TEST_F(RealtimeCommitPropertiesTest, BuildWithoutProgress) {
355368
RealtimeCommitProperties::Build(
356369
properties, std::optional<Snapshot>(MakeSnapshot(latest_properties)),
357370
/*realtime_ranges=*/{}, /*reset_all_realtime_progress=*/false,
358-
/*removed_realtime_partitions=*/{},
371+
/*removed_realtime_partitions=*/{}, *partition_computer_,
359372
/*file_system=*/nullptr,
360373
/*table_root=*/"", /*branch=*/"main"));
361374
ASSERT_EQ("value", inherited.at("custom"));
362375
ASSERT_EQ(latest_offsets_path, inherited.at(RealtimeCommitProperties::kOffsetsKey));
363376

364-
ASSERT_OK_AND_ASSIGN(Properties unchanged, RealtimeCommitProperties::Build(
365-
properties, /*latest_snapshot=*/std::nullopt,
366-
/*realtime_ranges=*/{},
367-
/*reset_all_realtime_progress=*/false,
368-
/*removed_realtime_partitions=*/{},
369-
/*file_system=*/nullptr,
370-
/*table_root=*/"", /*branch=*/"main"));
377+
ASSERT_OK_AND_ASSIGN(
378+
Properties unchanged,
379+
RealtimeCommitProperties::Build(properties, /*latest_snapshot=*/std::nullopt,
380+
/*realtime_ranges=*/{},
381+
/*reset_all_realtime_progress=*/false,
382+
/*removed_realtime_partitions=*/{}, *partition_computer_,
383+
/*file_system=*/nullptr,
384+
/*table_root=*/"", /*branch=*/"main"));
371385
ASSERT_EQ(properties, unchanged);
372386

373387
Properties properties_with_stale_offset = properties;
@@ -377,7 +391,7 @@ TEST_F(RealtimeCommitPropertiesTest, BuildWithoutProgress) {
377391
RealtimeCommitProperties::Build(
378392
properties_with_stale_offset, std::optional<Snapshot>(MakeSnapshot(latest_properties)),
379393
/*realtime_ranges=*/{}, /*reset_all_realtime_progress=*/true,
380-
/*removed_realtime_partitions=*/{},
394+
/*removed_realtime_partitions=*/{}, *partition_computer_,
381395
/*file_system=*/nullptr, /*table_root=*/"", /*branch=*/"main"));
382396
ASSERT_EQ("value", overwritten.at("custom"));
383397
ASSERT_EQ(0, overwritten.count(RealtimeCommitProperties::kOffsetsKey));
@@ -395,7 +409,7 @@ TEST_F(RealtimeCommitPropertiesTest, BuildRemovesOnlyOverwrittenPartitions) {
395409
RealtimeCommitProperties::Build(
396410
/*properties=*/{}, latest_snapshot, /*realtime_ranges=*/{},
397411
/*reset_all_realtime_progress=*/false, removed_partitions,
398-
file_system_, directory_->Str(), "main"));
412+
*partition_computer_, file_system_, directory_->Str(), "main"));
399413
ASSERT_OK_AND_ASSIGN(RealtimeOffsetMap actual,
400414
RealtimeCommitProperties::ReadOffsets(
401415
std::optional<Snapshot>(MakeSnapshot(properties)), file_system_));
@@ -418,7 +432,8 @@ TEST_F(RealtimeCommitPropertiesTest, BuildWritesMergedProgress) {
418432
RealtimeCommitProperties::Build(
419433
properties, std::optional<Snapshot>(MakeSnapshot(latest_properties)), ranges,
420434
/*reset_all_realtime_progress=*/false,
421-
/*removed_realtime_partitions=*/{}, file_system_, directory_->Str(), "main"));
435+
/*removed_realtime_partitions=*/{}, *partition_computer_, file_system_,
436+
directory_->Str(), "main"));
422437
ASSERT_EQ("value", merged.at("custom"));
423438
ASSERT_NE(latest_offsets_path, merged.at(RealtimeCommitProperties::kOffsetsKey));
424439

@@ -435,12 +450,12 @@ TEST_F(RealtimeCommitPropertiesTest, BuildWritesMergedProgress) {
435450
TEST_F(RealtimeCommitPropertiesTest, BuildRequiresFileSystem) {
436451
std::map<RealtimePartitionBucket, OffsetRange> ranges = {
437452
{RealtimePartitionBucket({{"dt", "2"}}, /*bucket=*/0), OffsetRange(0, 2)}};
438-
ASSERT_NOK_WITH_MSG(
439-
RealtimeCommitProperties::Build(
440-
/*properties=*/{}, /*latest_snapshot=*/std::nullopt, ranges,
441-
/*reset_all_realtime_progress=*/false,
442-
/*removed_realtime_partitions=*/{}, /*file_system=*/nullptr, directory_->Str(), "main"),
443-
"file system is null");
453+
ASSERT_NOK_WITH_MSG(RealtimeCommitProperties::Build(
454+
/*properties=*/{}, /*latest_snapshot=*/std::nullopt, ranges,
455+
/*reset_all_realtime_progress=*/false,
456+
/*removed_realtime_partitions=*/{}, *partition_computer_,
457+
/*file_system=*/nullptr, directory_->Str(), "main"),
458+
"file system is null");
444459
}
445460

446461
} // namespace paimon::test

0 commit comments

Comments
 (0)