[minor] Follow 3178, fix the flink metadata table compaction (#5175)
This commit is contained in:
@@ -207,8 +207,8 @@ public class TestStreamWriteOperatorCoordinator {
|
||||
assertThat(completedTimeline.lastInstant().get().getTimestamp(), is(HoodieTableMetadata.SOLO_COMMIT_TIMESTAMP));
|
||||
|
||||
// test metadata table compaction
|
||||
// write another 3 commits
|
||||
for (int i = 1; i < 4; i++) {
|
||||
// write another 4 commits
|
||||
for (int i = 1; i < 5; i++) {
|
||||
instant = mockWriteWithMetadata();
|
||||
metadataTableMetaClient.reloadActiveTimeline();
|
||||
completedTimeline = metadataTableMetaClient.getActiveTimeline().filterCompletedInstants();
|
||||
@@ -216,14 +216,14 @@ public class TestStreamWriteOperatorCoordinator {
|
||||
assertThat(completedTimeline.lastInstant().get().getTimestamp(), is(instant));
|
||||
}
|
||||
// the 5th commit triggers the compaction
|
||||
instant = mockWriteWithMetadata();
|
||||
mockWriteWithMetadata();
|
||||
metadataTableMetaClient.reloadActiveTimeline();
|
||||
completedTimeline = metadataTableMetaClient.getActiveTimeline().filterCompletedAndCompactionInstants();
|
||||
assertThat("One instant need to sync to metadata table", completedTimeline.getInstants().count(), is(6L));
|
||||
assertThat(completedTimeline.lastInstant().get().getTimestamp(), is(instant + "001"));
|
||||
assertThat(completedTimeline.lastInstant().get().getAction(), is(HoodieTimeline.COMMIT_ACTION));
|
||||
assertThat("One instant need to sync to metadata table", completedTimeline.getInstants().count(), is(7L));
|
||||
assertThat(completedTimeline.nthFromLastInstant(1).get().getTimestamp(), is(instant + "001"));
|
||||
assertThat(completedTimeline.nthFromLastInstant(1).get().getAction(), is(HoodieTimeline.COMMIT_ACTION));
|
||||
// write another 2 commits
|
||||
for (int i = 6; i < 8; i++) {
|
||||
for (int i = 7; i < 8; i++) {
|
||||
instant = mockWriteWithMetadata();
|
||||
metadataTableMetaClient.reloadActiveTimeline();
|
||||
completedTimeline = metadataTableMetaClient.getActiveTimeline().filterCompletedInstants();
|
||||
@@ -241,13 +241,15 @@ public class TestStreamWriteOperatorCoordinator {
|
||||
|
||||
// write another commit
|
||||
mockWriteWithMetadata();
|
||||
// write another commit to trigger compaction
|
||||
// write another commit
|
||||
instant = mockWriteWithMetadata();
|
||||
// write another commit to trigger compaction
|
||||
mockWriteWithMetadata();
|
||||
metadataTableMetaClient.reloadActiveTimeline();
|
||||
completedTimeline = metadataTableMetaClient.getActiveTimeline().filterCompletedAndCompactionInstants();
|
||||
assertThat("One instant need to sync to metadata table", completedTimeline.getInstants().count(), is(13L));
|
||||
assertThat(completedTimeline.lastInstant().get().getTimestamp(), is(instant + "001"));
|
||||
assertThat(completedTimeline.lastInstant().get().getAction(), is(HoodieTimeline.COMMIT_ACTION));
|
||||
assertThat("One instant need to sync to metadata table", completedTimeline.getInstants().count(), is(14L));
|
||||
assertThat(completedTimeline.nthFromLastInstant(1).get().getTimestamp(), is(instant + "001"));
|
||||
assertThat(completedTimeline.nthFromLastInstant(1).get().getAction(), is(HoodieTimeline.COMMIT_ACTION));
|
||||
}
|
||||
|
||||
// -------------------------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user