From 649fe6f3bba48b0c3355dabbda3f39197e79b0ac Mon Sep 17 00:00:00 2001 From: Lokesh Jain Date: Thu, 4 May 2023 15:32:41 +0530 Subject: [PATCH 1/3] HUDI-6170. Use correct zone id while calculating earliestTimeToRetain --- .../hudi/client/HoodieTimelineArchiver.java | 4 ++-- .../hudi/table/action/clean/CleanPlanner.java | 4 ++-- .../hudi/common/model/HoodieTimelineTimeZone.java | 15 ++++++++++++--- .../timeline/HoodieInstantTimeGenerator.java | 4 ++++ .../common/table/timeline/TestHoodieInstant.java | 9 +++++++++ 5 files changed, 29 insertions(+), 7 deletions(-) diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/HoodieTimelineArchiver.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/HoodieTimelineArchiver.java index 74e0a2565f827..42094e00d8283 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/HoodieTimelineArchiver.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/HoodieTimelineArchiver.java @@ -45,6 +45,7 @@ import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieArchivedTimeline; import org.apache.hudi.common.table.timeline.HoodieInstant; +import org.apache.hudi.common.table.timeline.HoodieInstantTimeGenerator; import org.apache.hudi.common.table.timeline.HoodieTimeline; import org.apache.hudi.common.table.timeline.TimelineMetadataUtils; import org.apache.hudi.common.table.timeline.TimelineUtils; @@ -77,7 +78,6 @@ import java.io.IOException; import java.text.ParseException; import java.time.Instant; -import java.time.ZoneId; import java.time.ZonedDateTime; import java.util.ArrayList; import java.util.Arrays; @@ -491,7 +491,7 @@ private Stream getCommitInstantsToArchive() throws IOException { String latestCommitToArchive = instantsToArchive.get(instantsToArchive.size() - 1).getTimestamp(); try { Instant latestCommitInstant = HoodieActiveTimeline.parseDateFromInstantTime(commitTimeline.lastInstant().get().getTimestamp()).toInstant(); - ZonedDateTime currentDateTime = ZonedDateTime.ofInstant(latestCommitInstant, ZoneId.systemDefault()); + ZonedDateTime currentDateTime = ZonedDateTime.ofInstant(latestCommitInstant, HoodieInstantTimeGenerator.getTimelineTimeZone().getZoneId()); String earliestTimeToRetain = HoodieActiveTimeline.formatDate(Date.from(currentDateTime.minusHours(config.getCleanerHoursRetained()).toInstant())); if (HoodieTimeline.compareTimestamps(latestCommitToArchive, GREATER_THAN_OR_EQUALS, earliestTimeToRetain)) { throw new HoodieIOException("Please align your archival configs based on cleaner configs. 'hoodie.keep.min.commits' : " diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/clean/CleanPlanner.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/clean/CleanPlanner.java index ae183b678e69a..7ac5e03eff059 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/clean/CleanPlanner.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/clean/CleanPlanner.java @@ -33,6 +33,7 @@ import org.apache.hudi.common.model.HoodieReplaceCommitMetadata; import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieInstant; +import org.apache.hudi.common.table.timeline.HoodieInstantTimeGenerator; import org.apache.hudi.common.table.timeline.HoodieTimeline; import org.apache.hudi.common.table.timeline.TimelineMetadataUtils; import org.apache.hudi.common.table.timeline.versioning.clean.CleanPlanV1MigrationHandler; @@ -53,7 +54,6 @@ import java.io.IOException; import java.io.Serializable; import java.time.Instant; -import java.time.ZoneId; import java.time.ZonedDateTime; import java.util.ArrayList; import java.util.Collections; @@ -510,7 +510,7 @@ public Option getEarliestCommitToRetain() { } } else if (config.getCleanerPolicy() == HoodieCleaningPolicy.KEEP_LATEST_BY_HOURS) { Instant instant = Instant.now(); - ZonedDateTime currentDateTime = ZonedDateTime.ofInstant(instant, ZoneId.systemDefault()); + ZonedDateTime currentDateTime = ZonedDateTime.ofInstant(instant, HoodieInstantTimeGenerator.getTimelineTimeZone().getZoneId()); String earliestTimeToRetain = HoodieActiveTimeline.formatDate(Date.from(currentDateTime.minusHours(hoursRetained).toInstant())); earliestCommitToRetain = Option.fromJavaOptional(commitTimeline.getInstantsAsStream().filter(i -> HoodieTimeline.compareTimestamps(i.getTimestamp(), HoodieTimeline.GREATER_THAN_OR_EQUALS, earliestTimeToRetain)).findFirst()); diff --git a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieTimelineTimeZone.java b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieTimelineTimeZone.java index 9b1c695d491ea..29e158812f485 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieTimelineTimeZone.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieTimelineTimeZone.java @@ -18,20 +18,29 @@ package org.apache.hudi.common.model; +import java.time.ZoneId; +import java.util.TimeZone; + /** * Hoodie TimelineZone. */ public enum HoodieTimelineTimeZone { - LOCAL("local"), - UTC("utc"); + LOCAL("local", ZoneId.systemDefault()), + UTC("utc", TimeZone.getTimeZone("UTC").toZoneId()); private final String timeZone; + private final ZoneId zoneId; - HoodieTimelineTimeZone(String timeZone) { + HoodieTimelineTimeZone(String timeZone, ZoneId zoneId) { this.timeZone = timeZone; + this.zoneId = zoneId; } public String getTimeZone() { return timeZone; } + + public ZoneId getZoneId() { + return zoneId; + } } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java index e839e73669e90..33773a6c917a3 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java @@ -144,4 +144,8 @@ public static boolean isValidInstantTime(String instantTime) { return false; } } + + public static HoodieTimelineTimeZone getTimelineTimeZone() { + return commitTimeZone; + } } diff --git a/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieInstant.java b/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieInstant.java index c4a1d00e90d06..39d4040c93abf 100644 --- a/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieInstant.java +++ b/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieInstant.java @@ -18,6 +18,7 @@ package org.apache.hudi.common.table.timeline; +import org.apache.hudi.common.model.HoodieTimelineTimeZone; import org.apache.hudi.common.testutils.HoodieCommonTestHarness; import org.apache.hudi.common.util.Option; import org.junit.jupiter.api.Test; @@ -27,6 +28,7 @@ import static org.apache.hudi.common.testutils.Assertions.assertStreamEquals; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; public class TestHoodieInstant extends HoodieCommonTestHarness { @@ -76,4 +78,11 @@ public void testCreateHoodieInstantByFileStatus() throws IOException { cleanMetaClient(); } } + + @Test + public void testHoodieTimelineTimeZone() { + for (HoodieTimelineTimeZone timeZone : HoodieTimelineTimeZone.values()) { + assertNotNull(timeZone.getZoneId()); + } + } } From 5e2f6a913cfb525f4f4e8a92c5009635a642478e Mon Sep 17 00:00:00 2001 From: Lokesh Jain Date: Mon, 8 May 2023 13:08:51 +0530 Subject: [PATCH 2/3] Address review comments --- .../org/apache/hudi/common/table/HoodieTableConfig.java | 4 +--- .../common/table/timeline/HoodieInstantTimeGenerator.java | 8 +++++--- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java b/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java index f4471e89a58f6..14214bed3998b 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java @@ -449,10 +449,8 @@ public static void create(FileSystem fs, Path metadataFolder, Properties propert // Use the default bootstrap index class. hoodieConfig.setDefaultValue(BOOTSTRAP_INDEX_CLASS_NAME, getDefaultBootstrapIndexClass(properties)); } - if (hoodieConfig.contains(TIMELINE_TIMEZONE)) { - HoodieInstantTimeGenerator.setCommitTimeZone(HoodieTimelineTimeZone.valueOf(hoodieConfig.getString(TIMELINE_TIMEZONE))); - } + HoodieInstantTimeGenerator.setCommitTimeZone(HoodieTimelineTimeZone.valueOf(hoodieConfig.getStringOrDefault(TIMELINE_TIMEZONE))); hoodieConfig.setDefaultValue(DROP_PARTITION_COLUMNS); storeProperties(hoodieConfig.getProps(), outputStream); diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java index 33773a6c917a3..bbd8c55b79819 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java @@ -19,6 +19,7 @@ package org.apache.hudi.common.table.timeline; import org.apache.hudi.common.model.HoodieTimelineTimeZone; +import org.apache.hudi.common.util.Option; import java.text.ParseException; import java.time.LocalDateTime; @@ -57,7 +58,7 @@ public class HoodieInstantTimeGenerator { // when performing comparisons such as LESS_THAN_OR_EQUAL_TO private static final String DEFAULT_MILLIS_EXT = "999"; - private static HoodieTimelineTimeZone commitTimeZone = HoodieTimelineTimeZone.LOCAL; + private static Option commitTimeZoneOpt = Option.empty(); /** * Returns next instant time that adds N milliseconds to the current time. @@ -66,6 +67,7 @@ public class HoodieInstantTimeGenerator { * @param milliseconds Milliseconds to add to current time while generating the new instant time */ public static String createNewInstantTime(long milliseconds) { + HoodieTimelineTimeZone commitTimeZone = commitTimeZoneOpt.get(); return lastInstantTime.updateAndGet((oldVal) -> { String newCommitTime; do { @@ -133,7 +135,7 @@ private static TemporalAccessor convertDateToTemporalAccessor(Date d) { } public static void setCommitTimeZone(HoodieTimelineTimeZone commitTimeZone) { - HoodieInstantTimeGenerator.commitTimeZone = commitTimeZone; + commitTimeZoneOpt = Option.of(commitTimeZone); } public static boolean isValidInstantTime(String instantTime) { @@ -146,6 +148,6 @@ public static boolean isValidInstantTime(String instantTime) { } public static HoodieTimelineTimeZone getTimelineTimeZone() { - return commitTimeZone; + return commitTimeZoneOpt.get(); } } From 1652e6ee6435b454dcbee662f7deaf7fb45883bf Mon Sep 17 00:00:00 2001 From: Lokesh Jain Date: Tue, 9 May 2023 12:43:55 +0530 Subject: [PATCH 3/3] Use metaClient table config --- .../org/apache/hudi/client/HoodieTimelineArchiver.java | 3 +-- .../apache/hudi/table/action/clean/CleanPlanner.java | 3 +-- .../apache/hudi/common/table/HoodieTableConfig.java | 8 +++++++- .../table/timeline/HoodieInstantTimeGenerator.java | 10 ++-------- 4 files changed, 11 insertions(+), 13 deletions(-) diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/HoodieTimelineArchiver.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/HoodieTimelineArchiver.java index 42094e00d8283..ac2c5042f74c7 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/HoodieTimelineArchiver.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/HoodieTimelineArchiver.java @@ -45,7 +45,6 @@ import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieArchivedTimeline; import org.apache.hudi.common.table.timeline.HoodieInstant; -import org.apache.hudi.common.table.timeline.HoodieInstantTimeGenerator; import org.apache.hudi.common.table.timeline.HoodieTimeline; import org.apache.hudi.common.table.timeline.TimelineMetadataUtils; import org.apache.hudi.common.table.timeline.TimelineUtils; @@ -491,7 +490,7 @@ private Stream getCommitInstantsToArchive() throws IOException { String latestCommitToArchive = instantsToArchive.get(instantsToArchive.size() - 1).getTimestamp(); try { Instant latestCommitInstant = HoodieActiveTimeline.parseDateFromInstantTime(commitTimeline.lastInstant().get().getTimestamp()).toInstant(); - ZonedDateTime currentDateTime = ZonedDateTime.ofInstant(latestCommitInstant, HoodieInstantTimeGenerator.getTimelineTimeZone().getZoneId()); + ZonedDateTime currentDateTime = ZonedDateTime.ofInstant(latestCommitInstant, metaClient.getTableConfig().getTimelineTimezone().getZoneId()); String earliestTimeToRetain = HoodieActiveTimeline.formatDate(Date.from(currentDateTime.minusHours(config.getCleanerHoursRetained()).toInstant())); if (HoodieTimeline.compareTimestamps(latestCommitToArchive, GREATER_THAN_OR_EQUALS, earliestTimeToRetain)) { throw new HoodieIOException("Please align your archival configs based on cleaner configs. 'hoodie.keep.min.commits' : " diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/clean/CleanPlanner.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/clean/CleanPlanner.java index 7ac5e03eff059..005e3f1dc344b 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/clean/CleanPlanner.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/clean/CleanPlanner.java @@ -33,7 +33,6 @@ import org.apache.hudi.common.model.HoodieReplaceCommitMetadata; import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieInstant; -import org.apache.hudi.common.table.timeline.HoodieInstantTimeGenerator; import org.apache.hudi.common.table.timeline.HoodieTimeline; import org.apache.hudi.common.table.timeline.TimelineMetadataUtils; import org.apache.hudi.common.table.timeline.versioning.clean.CleanPlanV1MigrationHandler; @@ -510,7 +509,7 @@ public Option getEarliestCommitToRetain() { } } else if (config.getCleanerPolicy() == HoodieCleaningPolicy.KEEP_LATEST_BY_HOURS) { Instant instant = Instant.now(); - ZonedDateTime currentDateTime = ZonedDateTime.ofInstant(instant, HoodieInstantTimeGenerator.getTimelineTimeZone().getZoneId()); + ZonedDateTime currentDateTime = ZonedDateTime.ofInstant(instant, hoodieTable.getMetaClient().getTableConfig().getTimelineTimezone().getZoneId()); String earliestTimeToRetain = HoodieActiveTimeline.formatDate(Date.from(currentDateTime.minusHours(hoursRetained).toInstant())); earliestCommitToRetain = Option.fromJavaOptional(commitTimeline.getInstantsAsStream().filter(i -> HoodieTimeline.compareTimestamps(i.getTimestamp(), HoodieTimeline.GREATER_THAN_OR_EQUALS, earliestTimeToRetain)).findFirst()); diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java b/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java index 14214bed3998b..22ba46a0956be 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableConfig.java @@ -449,8 +449,10 @@ public static void create(FileSystem fs, Path metadataFolder, Properties propert // Use the default bootstrap index class. hoodieConfig.setDefaultValue(BOOTSTRAP_INDEX_CLASS_NAME, getDefaultBootstrapIndexClass(properties)); } + if (hoodieConfig.contains(TIMELINE_TIMEZONE)) { + HoodieInstantTimeGenerator.setCommitTimeZone(HoodieTimelineTimeZone.valueOf(hoodieConfig.getString(TIMELINE_TIMEZONE))); + } - HoodieInstantTimeGenerator.setCommitTimeZone(HoodieTimelineTimeZone.valueOf(hoodieConfig.getStringOrDefault(TIMELINE_TIMEZONE))); hoodieConfig.setDefaultValue(DROP_PARTITION_COLUMNS); storeProperties(hoodieConfig.getProps(), outputStream); @@ -649,6 +651,10 @@ public String getKeyGeneratorClassName() { return getString(KEY_GENERATOR_CLASS_NAME); } + public HoodieTimelineTimeZone getTimelineTimezone() { + return HoodieTimelineTimeZone.valueOf(getStringOrDefault(TIMELINE_TIMEZONE)); + } + public String getHiveStylePartitioningEnable() { return getStringOrDefault(HIVE_STYLE_PARTITIONING_ENABLE); } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java index bbd8c55b79819..e839e73669e90 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieInstantTimeGenerator.java @@ -19,7 +19,6 @@ package org.apache.hudi.common.table.timeline; import org.apache.hudi.common.model.HoodieTimelineTimeZone; -import org.apache.hudi.common.util.Option; import java.text.ParseException; import java.time.LocalDateTime; @@ -58,7 +57,7 @@ public class HoodieInstantTimeGenerator { // when performing comparisons such as LESS_THAN_OR_EQUAL_TO private static final String DEFAULT_MILLIS_EXT = "999"; - private static Option commitTimeZoneOpt = Option.empty(); + private static HoodieTimelineTimeZone commitTimeZone = HoodieTimelineTimeZone.LOCAL; /** * Returns next instant time that adds N milliseconds to the current time. @@ -67,7 +66,6 @@ public class HoodieInstantTimeGenerator { * @param milliseconds Milliseconds to add to current time while generating the new instant time */ public static String createNewInstantTime(long milliseconds) { - HoodieTimelineTimeZone commitTimeZone = commitTimeZoneOpt.get(); return lastInstantTime.updateAndGet((oldVal) -> { String newCommitTime; do { @@ -135,7 +133,7 @@ private static TemporalAccessor convertDateToTemporalAccessor(Date d) { } public static void setCommitTimeZone(HoodieTimelineTimeZone commitTimeZone) { - commitTimeZoneOpt = Option.of(commitTimeZone); + HoodieInstantTimeGenerator.commitTimeZone = commitTimeZone; } public static boolean isValidInstantTime(String instantTime) { @@ -146,8 +144,4 @@ public static boolean isValidInstantTime(String instantTime) { return false; } } - - public static HoodieTimelineTimeZone getTimelineTimeZone() { - return commitTimeZoneOpt.get(); - } }