From 8c37d51f1b9f931211704a7b0196e715317e32db Mon Sep 17 00:00:00 2001 From: Dave Hagman Date: Mon, 18 Oct 2021 17:06:19 -0400 Subject: [PATCH 01/11] Update the instance format to include milliseconds --- .../apache/hudi/common/table/timeline/HoodieActiveTimeline.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java index e586815d3b974..a058cb81fa672 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java @@ -59,7 +59,7 @@ */ public class HoodieActiveTimeline extends HoodieDefaultTimeline { - public static final SimpleDateFormat COMMIT_FORMATTER = new SimpleDateFormat("yyyyMMddHHmmss"); + public static final SimpleDateFormat COMMIT_FORMATTER = new SimpleDateFormat("yyyyMMddHHmmssSSS"); public static final Set VALID_EXTENSIONS_IN_ACTIVE_TIMELINE = new HashSet<>(Arrays.asList( COMMIT_EXTENSION, INFLIGHT_COMMIT_EXTENSION, REQUESTED_COMMIT_EXTENSION, From 42f0681547fcb277c17452efde32de0643fc56db Mon Sep 17 00:00:00 2001 From: Dave Hagman Date: Tue, 19 Oct 2021 12:58:57 -0400 Subject: [PATCH 02/11] Abstract instant serde into active timeline methods --- .../org/apache/hudi/cli/utils/CommitUtil.java | 6 ++-- .../client/AbstractHoodieWriteClient.java | 2 +- .../hudi/client/HoodieFlinkWriteClient.java | 2 +- ...FlinkScheduleCompactionActionExecutor.java | 2 +- .../hudi/client/SparkRDDWriteClient.java | 4 +-- .../index/hbase/SparkHoodieHBaseIndex.java | 2 +- ...SparkScheduleCompactionActionExecutor.java | 2 +- .../table/timeline/HoodieActiveTimeline.java | 31 +++++++++++++++++-- .../apache/hudi/common/fs/TestFSUtils.java | 10 +++--- .../common/model/TestHoodieWriteStat.java | 4 +-- .../common/testutils/HoodieTestTable.java | 4 +-- .../org/apache/hudi/util/StreamerUtil.java | 10 +++--- .../spark/sql/hudi/HoodieSqlUtils.scala | 6 ++-- .../hudi/streaming/HoodieStreamSource.scala | 4 +-- .../hudi/functional/TestTimeTravelQuery.scala | 4 +-- .../functional/TestHDFSParquetImporter.java | 4 +-- 16 files changed, 62 insertions(+), 35 deletions(-) diff --git a/hudi-cli/src/main/java/org/apache/hudi/cli/utils/CommitUtil.java b/hudi-cli/src/main/java/org/apache/hudi/cli/utils/CommitUtil.java index 5a1c457b10ef1..1e067a3669c5d 100644 --- a/hudi-cli/src/main/java/org/apache/hudi/cli/utils/CommitUtil.java +++ b/hudi-cli/src/main/java/org/apache/hudi/cli/utils/CommitUtil.java @@ -51,7 +51,7 @@ public static long countNewRecords(HoodieTableMetaClient target, List co public static String getTimeDaysAgo(int numberOfDays) { Date date = Date.from(ZonedDateTime.now().minusDays(numberOfDays).toInstant()); - return HoodieActiveTimeline.COMMIT_FORMATTER.format(date); + return HoodieActiveTimeline.getInstantForDate(date); } /** @@ -61,8 +61,8 @@ public static String getTimeDaysAgo(int numberOfDays) { * b) hours: -1, returns 20200202010000 */ public static String addHours(String compactionCommitTime, int hours) throws ParseException { - Instant instant = HoodieActiveTimeline.COMMIT_FORMATTER.parse(compactionCommitTime).toInstant(); + Instant instant = HoodieActiveTimeline.parseDateFromInstantTime(compactionCommitTime).toInstant(); ZonedDateTime commitDateTime = ZonedDateTime.ofInstant(instant, ZoneId.systemDefault()); - return HoodieActiveTimeline.COMMIT_FORMATTER.format(Date.from(commitDateTime.plusHours(hours).toInstant())); + return HoodieActiveTimeline.getInstantForDate(Date.from(commitDateTime.plusHours(hours).toInstant())); } } diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/AbstractHoodieWriteClient.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/AbstractHoodieWriteClient.java index ec586a18034c3..ea04859a1dbf1 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/AbstractHoodieWriteClient.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/AbstractHoodieWriteClient.java @@ -232,7 +232,7 @@ void emitCommitMetrics(String instantTime, HoodieCommitMetadata metadata, String if (writeTimer != null) { long durationInMs = metrics.getDurationInMs(writeTimer.stop()); - metrics.updateCommitMetrics(HoodieActiveTimeline.COMMIT_FORMATTER.parse(instantTime).getTime(), durationInMs, + metrics.updateCommitMetrics(HoodieActiveTimeline.parseDateFromInstantTime(instantTime).getTime(), durationInMs, metadata, actionType); writeTimer = null; } diff --git a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java index 33878eb15693f..82455ce95bf99 100644 --- a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java +++ b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java @@ -367,7 +367,7 @@ public void completeCompaction( if (compactionTimer != null) { long durationInMs = metrics.getDurationInMs(compactionTimer.stop()); try { - metrics.updateCommitMetrics(HoodieActiveTimeline.COMMIT_FORMATTER.parse(compactionCommitTime).getTime(), + metrics.updateCommitMetrics(HoodieActiveTimeline.parseDateFromInstantTime(compactionCommitTime).getTime(), durationInMs, metadata, HoodieActiveTimeline.COMPACTION_ACTION); } catch (ParseException e) { throw new HoodieCommitException("Commit time is not of valid format. Failed to commit compaction " diff --git a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/compact/FlinkScheduleCompactionActionExecutor.java b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/compact/FlinkScheduleCompactionActionExecutor.java index 4143944bbebc8..25162bf260fd6 100644 --- a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/compact/FlinkScheduleCompactionActionExecutor.java +++ b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/compact/FlinkScheduleCompactionActionExecutor.java @@ -147,7 +147,7 @@ public boolean needCompact(CompactionTriggerStrategy compactionTriggerStrategy) public Long parsedToSeconds(String time) { long timestamp; try { - timestamp = HoodieActiveTimeline.COMMIT_FORMATTER.parse(time).getTime() / 1000; + timestamp = HoodieActiveTimeline.parseDateFromInstantTime(time).getTime() / 1000; } catch (ParseException e) { throw new HoodieCompactionException(e.getMessage(), e); } diff --git a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/SparkRDDWriteClient.java b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/SparkRDDWriteClient.java index dd9f43d16a901..92269e653676c 100644 --- a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/SparkRDDWriteClient.java +++ b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/SparkRDDWriteClient.java @@ -313,7 +313,7 @@ protected void completeCompaction(HoodieCommitMetadata metadata, JavaRDD VALID_EXTENSIONS_IN_ACTIVE_TIMELINE = new HashSet<>(Arrays.asList( COMMIT_EXTENSION, INFLIGHT_COMMIT_EXTENSION, REQUESTED_COMMIT_EXTENSION, @@ -74,6 +77,26 @@ public class HoodieActiveTimeline extends HoodieDefaultTimeline { protected HoodieTableMetaClient metaClient; private static AtomicReference lastInstantTime = new AtomicReference<>(String.valueOf(Integer.MIN_VALUE)); + /** + * Parses the given instant ID to return a date instance. + * @param instant The instant ID + * @return A date + * @throws ParseException If the instant ID is malformed + */ + public static Date parseDateFromInstantTime(String instant) throws ParseException { + // Enables backwards compatibility with non-millisecond granularity instants + if (isMillisecondGranularity(instant)) { + return COMMIT_FORMATTER.parse(instant); + } else { + // Add milliseconds to the instant in order to parse successfully + return COMMIT_FORMATTER.parse(instant + "000"); + } + } + + public static String getInstantForDate(Date instantDate) { + return COMMIT_FORMATTER.format(instantDate); + } + /** * Returns next instant time in the {@link #COMMIT_FORMATTER} format. * Ensures each instant time is atleast 1 second apart since we create instant times at second granularity @@ -90,7 +113,7 @@ public static String createNewInstantTime(long milliseconds) { return lastInstantTime.updateAndGet((oldVal) -> { String newCommitTime; do { - newCommitTime = HoodieActiveTimeline.COMMIT_FORMATTER.format(new Date(System.currentTimeMillis() + milliseconds)); + newCommitTime = getInstantForDate(new Date(System.currentTimeMillis() + milliseconds)); } while (HoodieTimeline.compareTimestamps(newCommitTime, LESSER_THAN_OR_EQUALS, oldVal)); return newCommitTime; }); @@ -193,6 +216,10 @@ public void deleteCompactionRequested(HoodieInstant instant) { deleteInstantFile(instant); } + private static boolean isMillisecondGranularity(String instant) { + return instant.length() == INSTANT_ID_LENGTH; + } + private void deleteInstantFileIfExists(HoodieInstant instant) { LOG.info("Deleting instant " + instant); Path inFlightCommitFilePath = new Path(metaClient.getMetaPath(), instant.getFileName()); diff --git a/hudi-common/src/test/java/org/apache/hudi/common/fs/TestFSUtils.java b/hudi-common/src/test/java/org/apache/hudi/common/fs/TestFSUtils.java index ef8b09b51e440..1e46cd5c35fb2 100644 --- a/hudi-common/src/test/java/org/apache/hudi/common/fs/TestFSUtils.java +++ b/hudi-common/src/test/java/org/apache/hudi/common/fs/TestFSUtils.java @@ -23,6 +23,7 @@ import org.apache.hudi.common.model.HoodieLogFile; import org.apache.hudi.common.table.HoodieTableConfig; import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.testutils.HoodieCommonTestHarness; import org.apache.hudi.common.testutils.HoodieTestUtils; import org.apache.hudi.exception.HoodieException; @@ -51,7 +52,6 @@ import java.util.stream.Stream; import static org.apache.hudi.common.model.HoodieFileFormat.HOODIE_LOG; -import static org.apache.hudi.common.table.timeline.HoodieActiveTimeline.COMMIT_FORMATTER; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNull; @@ -79,14 +79,14 @@ public void setUp() throws IOException { @Test public void testMakeDataFileName() { - String instantTime = COMMIT_FORMATTER.format(new Date()); + String instantTime = HoodieActiveTimeline.getInstantForDate(new Date()); String fileName = UUID.randomUUID().toString(); assertEquals(FSUtils.makeDataFileName(instantTime, TEST_WRITE_TOKEN, fileName), fileName + "_" + TEST_WRITE_TOKEN + "_" + instantTime + BASE_FILE_EXTENSION); } @Test public void testMaskFileName() { - String instantTime = COMMIT_FORMATTER.format(new Date()); + String instantTime = HoodieActiveTimeline.getInstantForDate(new Date()); int taskPartitionId = 2; assertEquals(FSUtils.maskWithoutFileId(instantTime, taskPartitionId), "*_" + taskPartitionId + "_" + instantTime + BASE_FILE_EXTENSION); } @@ -154,7 +154,7 @@ public void testProcessFiles() throws Exception { @Test public void testGetCommitTime() { - String instantTime = COMMIT_FORMATTER.format(new Date()); + String instantTime = HoodieActiveTimeline.getInstantForDate(new Date()); String fileName = UUID.randomUUID().toString(); String fullFileName = FSUtils.makeDataFileName(instantTime, TEST_WRITE_TOKEN, fileName); assertEquals(instantTime, FSUtils.getCommitTime(fullFileName)); @@ -165,7 +165,7 @@ public void testGetCommitTime() { @Test public void testGetFileNameWithoutMeta() { - String instantTime = COMMIT_FORMATTER.format(new Date()); + String instantTime = HoodieActiveTimeline.getInstantForDate(new Date()); String fileName = UUID.randomUUID().toString(); String fullFileName = FSUtils.makeDataFileName(instantTime, TEST_WRITE_TOKEN, fileName); assertEquals(fileName, FSUtils.getFileId(fullFileName)); diff --git a/hudi-common/src/test/java/org/apache/hudi/common/model/TestHoodieWriteStat.java b/hudi-common/src/test/java/org/apache/hudi/common/model/TestHoodieWriteStat.java index 7136ce7d372bb..89588408547b2 100644 --- a/hudi-common/src/test/java/org/apache/hudi/common/model/TestHoodieWriteStat.java +++ b/hudi-common/src/test/java/org/apache/hudi/common/model/TestHoodieWriteStat.java @@ -21,12 +21,12 @@ import org.apache.hudi.common.fs.FSUtils; import org.apache.hadoop.fs.Path; +import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.junit.jupiter.api.Test; import java.util.Date; import java.util.UUID; -import static org.apache.hudi.common.table.timeline.HoodieActiveTimeline.COMMIT_FORMATTER; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNull; @@ -37,7 +37,7 @@ public class TestHoodieWriteStat { @Test public void testSetPaths() { - String instantTime = COMMIT_FORMATTER.format(new Date()); + String instantTime = HoodieActiveTimeline.getInstantForDate(new Date()); String basePathString = "/data/tables/some-hoodie-table"; String partitionPathString = "2017/12/31"; String fileName = UUID.randomUUID().toString(); diff --git a/hudi-common/src/test/java/org/apache/hudi/common/testutils/HoodieTestTable.java b/hudi-common/src/test/java/org/apache/hudi/common/testutils/HoodieTestTable.java index 2a829b596a685..f775ac8bc9475 100644 --- a/hudi-common/src/test/java/org/apache/hudi/common/testutils/HoodieTestTable.java +++ b/hudi-common/src/test/java/org/apache/hudi/common/testutils/HoodieTestTable.java @@ -44,6 +44,7 @@ import org.apache.hudi.common.model.WriteOperationType; import org.apache.hudi.common.table.HoodieTableConfig; import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieInstant; import org.apache.hudi.common.table.timeline.HoodieTimeline; import org.apache.hudi.common.table.timeline.TimelineMetadataUtils; @@ -84,7 +85,6 @@ import static org.apache.hudi.common.model.WriteOperationType.CLUSTER; import static org.apache.hudi.common.model.WriteOperationType.COMPACT; import static org.apache.hudi.common.model.WriteOperationType.UPSERT; -import static org.apache.hudi.common.table.timeline.HoodieActiveTimeline.COMMIT_FORMATTER; import static org.apache.hudi.common.table.timeline.HoodieTimeline.CLEAN_ACTION; import static org.apache.hudi.common.table.timeline.HoodieTimeline.REPLACE_COMMIT_ACTION; import static org.apache.hudi.common.testutils.FileCreateUtils.baseFileName; @@ -147,7 +147,7 @@ public static String makeNewCommitTime() { } public static String makeNewCommitTime(Instant dateTime) { - return COMMIT_FORMATTER.format(Date.from(dateTime)); + return HoodieActiveTimeline.getInstantForDate(Date.from(dateTime)); } public static List makeIncrementalCommitTimes(int num) { diff --git a/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java b/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java index 835bb49b42f38..58971a7d0ae6c 100644 --- a/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java +++ b/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java @@ -407,12 +407,12 @@ public static HoodieFlinkWriteClient createWriteClient(Configuration conf) throw */ public static String medianInstantTime(String highVal, String lowVal) { try { - long high = HoodieActiveTimeline.COMMIT_FORMATTER.parse(highVal).getTime(); - long low = HoodieActiveTimeline.COMMIT_FORMATTER.parse(lowVal).getTime(); + long high = HoodieActiveTimeline.parseDateFromInstantTime(highVal).getTime(); + long low = HoodieActiveTimeline.parseDateFromInstantTime(lowVal).getTime(); ValidationUtils.checkArgument(high > low, "Instant [" + highVal + "] should have newer timestamp than instant [" + lowVal + "]"); long median = low + (high - low) / 2; - return HoodieActiveTimeline.COMMIT_FORMATTER.format(new Date(median)); + return HoodieActiveTimeline.getInstantForDate(new Date(median)); } catch (ParseException e) { throw new HoodieException("Get median instant time with interval [" + lowVal + ", " + highVal + "] error", e); } @@ -423,8 +423,8 @@ public static String medianInstantTime(String highVal, String lowVal) { */ public static long instantTimeDiffSeconds(String newInstantTime, String oldInstantTime) { try { - long newTimestamp = HoodieActiveTimeline.COMMIT_FORMATTER.parse(newInstantTime).getTime(); - long oldTimestamp = HoodieActiveTimeline.COMMIT_FORMATTER.parse(oldInstantTime).getTime(); + long newTimestamp = HoodieActiveTimeline.parseDateFromInstantTime(newInstantTime).getTime(); + long oldTimestamp = HoodieActiveTimeline.parseDateFromInstantTime(oldInstantTime).getTime(); return (newTimestamp - oldTimestamp) / 1000; } catch (ParseException e) { throw new HoodieException("Get instant time diff with interval [" + oldInstantTime + ", " + newInstantTime + "] error", e); diff --git a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/HoodieSqlUtils.scala b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/HoodieSqlUtils.scala index 318577b81410f..4d4053a35f45c 100644 --- a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/HoodieSqlUtils.scala +++ b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/HoodieSqlUtils.scala @@ -254,12 +254,12 @@ object HoodieSqlUtils extends SparkAdapterSupport { */ def formatQueryInstant(queryInstant: String): String = { if (queryInstant.length == 19) { // for yyyy-MM-dd HH:mm:ss - HoodieActiveTimeline.COMMIT_FORMATTER.format(defaultDateTimeFormat.parse(queryInstant)) + HoodieActiveTimeline.getInstantForDate(defaultDateTimeFormat.parse(queryInstant)) } else if (queryInstant.length == 14) { // for yyyyMMddHHmmss - HoodieActiveTimeline.COMMIT_FORMATTER.parse(queryInstant) // validate the format + HoodieActiveTimeline.parseDateFromInstantTime(queryInstant) // validate the format queryInstant } else if (queryInstant.length == 10) { // for yyyy-MM-dd - HoodieActiveTimeline.COMMIT_FORMATTER.format(defaultDateFormat.parse(queryInstant)) + HoodieActiveTimeline.getInstantForDate(defaultDateFormat.parse(queryInstant)) } else { throw new IllegalArgumentException(s"Unsupported query instant time format: $queryInstant," + s"Supported time format are: 'yyyy-MM-dd: HH:mm:ss' or 'yyyy-MM-dd' or 'yyyyMMddHHmmss'") diff --git a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSource.scala b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSource.scala index 0482e74884926..900fff14e0f6a 100644 --- a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSource.scala +++ b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSource.scala @@ -179,10 +179,10 @@ class HoodieStreamSource( startOffset match { case INIT_OFFSET => startOffset.commitTime case HoodieSourceOffset(commitTime) => - val time = HoodieActiveTimeline.COMMIT_FORMATTER.parse(commitTime).getTime + val time = HoodieActiveTimeline.parseDateFromInstantTime(commitTime).getTime // As we consume the data between (start, end], start is not included, // so we +1s to the start commit time here. - HoodieActiveTimeline.COMMIT_FORMATTER.format(new Date(time + 1000)) + HoodieActiveTimeline.getInstantForDate(new Date(time + 1000)) case _=> throw new IllegalStateException("UnKnow offset type.") } } diff --git a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestTimeTravelQuery.scala b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestTimeTravelQuery.scala index bb102a4cd912e..9064800525dc1 100644 --- a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestTimeTravelQuery.scala +++ b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestTimeTravelQuery.scala @@ -217,13 +217,13 @@ class TestTimeTravelQuery extends HoodieClientTestBase { } private def defaultDateTimeFormat(queryInstant: String): String = { - val date = HoodieActiveTimeline.COMMIT_FORMATTER.parse(queryInstant) + val date = HoodieActiveTimeline.parseDateFromInstantTime(queryInstant) val format = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss") format.format(date) } private def defaultDateFormat(queryInstant: String): String = { - val date = HoodieActiveTimeline.COMMIT_FORMATTER.parse(queryInstant) + val date = HoodieActiveTimeline.parseDateFromInstantTime(queryInstant) val format = new SimpleDateFormat("yyyy-MM-dd") format.format(date) } diff --git a/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHDFSParquetImporter.java b/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHDFSParquetImporter.java index 6d0141e407b88..3ac490bf9163e 100644 --- a/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHDFSParquetImporter.java +++ b/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHDFSParquetImporter.java @@ -230,7 +230,7 @@ public void testImportWithUpsert() throws IOException, ParseException { public List createInsertRecords(Path srcFolder) throws ParseException, IOException { Path srcFile = new Path(srcFolder.toString(), "file1.parquet"); - long startTime = HoodieActiveTimeline.COMMIT_FORMATTER.parse("20170203000000").getTime() / 1000; + long startTime = HoodieActiveTimeline.parseDateFromInstantTime("20170203000000").getTime() / 1000; List records = new ArrayList(); for (long recordNum = 0; recordNum < 96; recordNum++) { records.add(HoodieTestDataGenerator.generateGenericRecord(Long.toString(recordNum), "0", "rider-" + recordNum, @@ -247,7 +247,7 @@ public List createInsertRecords(Path srcFolder) throws ParseExcep public List createUpsertRecords(Path srcFolder) throws ParseException, IOException { Path srcFile = new Path(srcFolder.toString(), "file1.parquet"); - long startTime = HoodieActiveTimeline.COMMIT_FORMATTER.parse("20170203000000").getTime() / 1000; + long startTime = HoodieActiveTimeline.parseDateFromInstantTime("20170203000000").getTime() / 1000; List records = new ArrayList(); // 10 for update for (long recordNum = 0; recordNum < 11; recordNum++) { From 4371002aeba0756c6e2398b68915eda432759b7a Mon Sep 17 00:00:00 2001 From: Dave Hagman Date: Wed, 20 Oct 2021 11:43:50 -0400 Subject: [PATCH 03/11] Add test to verify parsing and instant date math works as expected --- .../timeline/TestHoodieActiveTimeline.java | 23 +++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieActiveTimeline.java b/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieActiveTimeline.java index 5c4c911e1576d..f802cfc20bb79 100755 --- a/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieActiveTimeline.java +++ b/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieActiveTimeline.java @@ -31,8 +31,10 @@ import org.junit.jupiter.api.Test; import java.io.IOException; +import java.text.ParseException; import java.util.ArrayList; import java.util.Collections; +import java.util.Date; import java.util.HashSet; import java.util.List; import java.util.Random; @@ -428,6 +430,27 @@ public void testReplaceActionsTimeline() { assertEquals(HoodieTimeline.REPLACE_COMMIT_ACTION, validReplaceInstants.get(0).getAction()); } + @Test + public void testMillisGranularityInstantDateParsing() throws ParseException { + // Old second granularity instant ID + String secondGranularityInstant = "20210101120101"; + Date secondGranularityDate = HoodieActiveTimeline.parseDateFromInstantTime(secondGranularityInstant); + // New ms granularity instant ID + String msGranularityInstant = secondGranularityInstant + "009"; + Date msGranularityDate = HoodieActiveTimeline.parseDateFromInstantTime(msGranularityInstant); + assertEquals(0, secondGranularityDate.getTime() % 1000, "Expected the ms part to be 0"); + assertEquals(9, msGranularityDate.getTime() % 1000, "Expected the ms part to be 9"); + + // Ensure that any date math which expects second granularity still works + String laterDateInstant = "20210101120111"; // + 10 seconds from original instant + assertEquals( + 10, + HoodieActiveTimeline.parseDateFromInstantTime(laterDateInstant).getTime() / 1000 + - HoodieActiveTimeline.parseDateFromInstantTime(secondGranularityInstant).getTime() / 1000, + "Expected the difference between later instant and previous instant to be 10 seconds" + ); + } + /** * Returns an exhaustive list of all possible HoodieInstant. * @return list of HoodieInstant From 719f111bc62e405066bd76168e896a8f9e78d1f7 Mon Sep 17 00:00:00 2001 From: Dave Hagman Date: Wed, 20 Oct 2021 15:25:01 -0400 Subject: [PATCH 04/11] Update returned test value to be ms granularity --- .../src/test/java/org/apache/hudi/utils/TestStreamerUtil.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java b/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java index af40a8dd82e60..263988346a3a9 100644 --- a/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java +++ b/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java @@ -81,7 +81,7 @@ void testMedianInstantTime() { String higher = "20210705125921"; String lower = "20210705125806"; String median1 = StreamerUtil.medianInstantTime(higher, lower); - assertThat(median1, is("20210705125843")); + assertThat(median1, is("20210705125843500")); // test symmetry assertThrows(IllegalArgumentException.class, () -> StreamerUtil.medianInstantTime(lower, higher), From 7d16138887d87c5dc9f28383c94e4407c0a4e1cf Mon Sep 17 00:00:00 2001 From: Dave Hagman Date: Mon, 25 Oct 2021 11:31:14 -0400 Subject: [PATCH 05/11] Add method to timeline to correctly parse incoming date-times for instant queries --- .../table/timeline/HoodieActiveTimeline.java | 24 +++++++++++++++++++ .../spark/sql/hudi/HoodieSqlUtils.scala | 8 +++---- 2 files changed, 28 insertions(+), 4 deletions(-) diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java index b8770c14fc74c..e0bad763eae74 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java @@ -64,6 +64,13 @@ public class HoodieActiveTimeline extends HoodieDefaultTimeline { private static final int INSTANT_ID_LENGTH = COMMIT_FORMAT.length(); private static final SimpleDateFormat COMMIT_FORMATTER = new SimpleDateFormat(COMMIT_FORMAT); + private static final String MILLIS_GRANULARITY_DATE_FORMAT = "yyyy-MM-dd HH:mm:ss:SSS"; + private static final SimpleDateFormat MS_GRANULARITY_DATE_FORMATTER = new SimpleDateFormat(MILLIS_GRANULARITY_DATE_FORMAT); + + // The default number of milliseconds that we add if they are not present + // We prefer the max timestamp as it mimics the current behavior with second granularity + private static final String DEFAULT_MILLIS_EXT = "999"; + public static final Set VALID_EXTENSIONS_IN_ACTIVE_TIMELINE = new HashSet<>(Arrays.asList( COMMIT_EXTENSION, INFLIGHT_COMMIT_EXTENSION, REQUESTED_COMMIT_EXTENSION, DELTA_COMMIT_EXTENSION, INFLIGHT_DELTA_COMMIT_EXTENSION, REQUESTED_DELTA_COMMIT_EXTENSION, @@ -97,6 +104,23 @@ public static String getInstantForDate(Date instantDate) { return COMMIT_FORMATTER.format(instantDate); } + /** + * Creates an instant string given a valid date-time string. + * @param dateString A date-time string in the format yyyy-MM-dd HH:mm:ss[:SSS] + * @return A timeline instant + * @throws ParseException If we cannot parse the date string + */ + public static String getInstantForDateString(String dateString) throws ParseException { + try { + return getInstantForDate(MS_GRANULARITY_DATE_FORMATTER.parse(dateString)); + } catch (ParseException e) { + // Attempt to add the milliseconds in order to complete parsing + return getInstantForDate(MS_GRANULARITY_DATE_FORMATTER.parse( + String.format("%s:%s", dateString, DEFAULT_MILLIS_EXT) + )); + } + } + /** * Returns next instant time in the {@link #COMMIT_FORMATTER} format. * Ensures each instant time is atleast 1 second apart since we create instant times at second granularity diff --git a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/HoodieSqlUtils.scala b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/HoodieSqlUtils.scala index 4d4053a35f45c..923366bdd4537 100644 --- a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/HoodieSqlUtils.scala +++ b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/HoodieSqlUtils.scala @@ -47,7 +47,6 @@ import java.text.SimpleDateFormat import scala.collection.immutable.Map object HoodieSqlUtils extends SparkAdapterSupport { - private val defaultDateTimeFormat = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss") private val defaultDateFormat = new SimpleDateFormat("yyyy-MM-dd") def isHoodieTable(table: CatalogTable): Boolean = { @@ -253,9 +252,10 @@ object HoodieSqlUtils extends SparkAdapterSupport { * 3、yyyyMMddHHmmss */ def formatQueryInstant(queryInstant: String): String = { - if (queryInstant.length == 19) { // for yyyy-MM-dd HH:mm:ss - HoodieActiveTimeline.getInstantForDate(defaultDateTimeFormat.parse(queryInstant)) - } else if (queryInstant.length == 14) { // for yyyyMMddHHmmss + val instantLength = queryInstant.length + if (instantLength == 19 || instantLength == 23) { // for yyyy-MM-dd HH:mm:ss[:SSS] + HoodieActiveTimeline.getInstantForDateString(queryInstant) + } else if (instantLength == 14 || instantLength == 17) { // for yyyyMMddHHmmss[SSS] HoodieActiveTimeline.parseDateFromInstantTime(queryInstant) // validate the format queryInstant } else if (queryInstant.length == 10) { // for yyyy-MM-dd From 50cc95c8801b0770eaf137b9f691b6b570c23d52 Mon Sep 17 00:00:00 2001 From: Dave Hagman Date: Tue, 26 Oct 2021 13:26:35 -0400 Subject: [PATCH 06/11] Use default millis --- .../apache/hudi/common/table/timeline/HoodieActiveTimeline.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java index e0bad763eae74..3f1f627fe6c4e 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java @@ -96,7 +96,7 @@ public static Date parseDateFromInstantTime(String instant) throws ParseExceptio return COMMIT_FORMATTER.parse(instant); } else { // Add milliseconds to the instant in order to parse successfully - return COMMIT_FORMATTER.parse(instant + "000"); + return COMMIT_FORMATTER.parse(instant + DEFAULT_MILLIS_EXT); } } From 1209fe7ad83387b945d789346b068fac67aef759 Mon Sep 17 00:00:00 2001 From: Dave Hagman Date: Tue, 26 Oct 2021 13:29:07 -0400 Subject: [PATCH 07/11] Make var naming more consistent. Better comments --- .../table/timeline/HoodieActiveTimeline.java | 15 ++++++++------- 1 file changed, 8 insertions(+), 7 deletions(-) diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java index 3f1f627fe6c4e..09c8d292f50aa 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java @@ -60,15 +60,16 @@ */ public class HoodieActiveTimeline extends HoodieDefaultTimeline { - private static final String COMMIT_FORMAT = "yyyyMMddHHmmssSSS"; - private static final int INSTANT_ID_LENGTH = COMMIT_FORMAT.length(); - private static final SimpleDateFormat COMMIT_FORMATTER = new SimpleDateFormat(COMMIT_FORMAT); + private static final String MILLIS_COMMIT_FORMAT = "yyyyMMddHHmmssSSS"; + private static final int MILLIS_INSTANT_ID_LENGTH = MILLIS_COMMIT_FORMAT.length(); + private static final SimpleDateFormat COMMIT_FORMATTER = new SimpleDateFormat(MILLIS_COMMIT_FORMAT); private static final String MILLIS_GRANULARITY_DATE_FORMAT = "yyyy-MM-dd HH:mm:ss:SSS"; - private static final SimpleDateFormat MS_GRANULARITY_DATE_FORMATTER = new SimpleDateFormat(MILLIS_GRANULARITY_DATE_FORMAT); + private static final SimpleDateFormat MILLIS_GRANULARITY_DATE_FORMATTER = new SimpleDateFormat(MILLIS_GRANULARITY_DATE_FORMAT); // The default number of milliseconds that we add if they are not present // We prefer the max timestamp as it mimics the current behavior with second granularity + // when performing comparisons such as LESS_THAN_OR_EQUAL_TO private static final String DEFAULT_MILLIS_EXT = "999"; public static final Set VALID_EXTENSIONS_IN_ACTIVE_TIMELINE = new HashSet<>(Arrays.asList( @@ -112,10 +113,10 @@ public static String getInstantForDate(Date instantDate) { */ public static String getInstantForDateString(String dateString) throws ParseException { try { - return getInstantForDate(MS_GRANULARITY_DATE_FORMATTER.parse(dateString)); + return getInstantForDate(MILLIS_GRANULARITY_DATE_FORMATTER.parse(dateString)); } catch (ParseException e) { // Attempt to add the milliseconds in order to complete parsing - return getInstantForDate(MS_GRANULARITY_DATE_FORMATTER.parse( + return getInstantForDate(MILLIS_GRANULARITY_DATE_FORMATTER.parse( String.format("%s:%s", dateString, DEFAULT_MILLIS_EXT) )); } @@ -241,7 +242,7 @@ public void deleteCompactionRequested(HoodieInstant instant) { } private static boolean isMillisecondGranularity(String instant) { - return instant.length() == INSTANT_ID_LENGTH; + return instant.length() == MILLIS_INSTANT_ID_LENGTH; } private void deleteInstantFileIfExists(HoodieInstant instant) { From d4a8853e50a322d928d30fe01bc8d96a4f9814cb Mon Sep 17 00:00:00 2001 From: Dave Hagman Date: Wed, 27 Oct 2021 09:40:48 -0400 Subject: [PATCH 08/11] Fixed test to line up with new default millis on instants --- .../common/table/timeline/TestHoodieActiveTimeline.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieActiveTimeline.java b/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieActiveTimeline.java index f802cfc20bb79..1e37d5edd0577 100755 --- a/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieActiveTimeline.java +++ b/hudi-common/src/test/java/org/apache/hudi/common/table/timeline/TestHoodieActiveTimeline.java @@ -434,11 +434,11 @@ public void testReplaceActionsTimeline() { public void testMillisGranularityInstantDateParsing() throws ParseException { // Old second granularity instant ID String secondGranularityInstant = "20210101120101"; - Date secondGranularityDate = HoodieActiveTimeline.parseDateFromInstantTime(secondGranularityInstant); + Date defaultMsGranularityDate = HoodieActiveTimeline.parseDateFromInstantTime(secondGranularityInstant); // New ms granularity instant ID - String msGranularityInstant = secondGranularityInstant + "009"; - Date msGranularityDate = HoodieActiveTimeline.parseDateFromInstantTime(msGranularityInstant); - assertEquals(0, secondGranularityDate.getTime() % 1000, "Expected the ms part to be 0"); + String specificMsGranularityInstant = secondGranularityInstant + "009"; + Date msGranularityDate = HoodieActiveTimeline.parseDateFromInstantTime(specificMsGranularityInstant); + assertEquals(999, defaultMsGranularityDate.getTime() % 1000, "Expected the ms part to be 999"); assertEquals(9, msGranularityDate.getTime() % 1000, "Expected the ms part to be 9"); // Ensure that any date math which expects second granularity still works From 52b7ea1f9d000003d2a31dd88fc41a35e121827b Mon Sep 17 00:00:00 2001 From: Dave Hagman Date: Wed, 27 Oct 2021 13:56:29 -0400 Subject: [PATCH 09/11] Update test value (which is actually correct). Add ability to override the default number of milliseconds added to second granularity instants --- .../table/timeline/HoodieActiveTimeline.java | 31 ++++++++++++++----- .../apache/hudi/utils/TestStreamerUtil.java | 3 +- 2 files changed, 25 insertions(+), 9 deletions(-) diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java index 09c8d292f50aa..6516917d6b1eb 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java @@ -67,10 +67,11 @@ public class HoodieActiveTimeline extends HoodieDefaultTimeline { private static final String MILLIS_GRANULARITY_DATE_FORMAT = "yyyy-MM-dd HH:mm:ss:SSS"; private static final SimpleDateFormat MILLIS_GRANULARITY_DATE_FORMATTER = new SimpleDateFormat(MILLIS_GRANULARITY_DATE_FORMAT); - // The default number of milliseconds that we add if they are not present - // We prefer the max timestamp as it mimics the current behavior with second granularity - // when performing comparisons such as LESS_THAN_OR_EQUAL_TO - private static final String DEFAULT_MILLIS_EXT = "999"; + // Millisecond extensions available when converting from second -> millisecond granularity + // You may prefer the max timestamp as it mimics the current behavior with second granularity + // when performing comparisons such as LESS_THAN_OR_EQUAL_TO but there are times where the min is desirable. + private static final String MIN_MILLIS_INSTANT_EXT = "000"; + private static final String MAX_MILLIS_INSTANT_EXT = "999"; public static final Set VALID_EXTENSIONS_IN_ACTIVE_TIMELINE = new HashSet<>(Arrays.asList( COMMIT_EXTENSION, INFLIGHT_COMMIT_EXTENSION, REQUESTED_COMMIT_EXTENSION, @@ -88,19 +89,26 @@ public class HoodieActiveTimeline extends HoodieDefaultTimeline { /** * Parses the given instant ID to return a date instance. * @param instant The instant ID + * @param preferMaxDefaultMillis If true and dateString is missing milliseconds then + * append the max number of milliseconds. If false it will be min * @return A date * @throws ParseException If the instant ID is malformed */ - public static Date parseDateFromInstantTime(String instant) throws ParseException { + public static Date parseDateFromInstantTime(String instant, boolean preferMaxDefaultMillis) throws ParseException { // Enables backwards compatibility with non-millisecond granularity instants if (isMillisecondGranularity(instant)) { return COMMIT_FORMATTER.parse(instant); } else { // Add milliseconds to the instant in order to parse successfully - return COMMIT_FORMATTER.parse(instant + DEFAULT_MILLIS_EXT); + String ext = preferMaxDefaultMillis ? MAX_MILLIS_INSTANT_EXT : MIN_MILLIS_INSTANT_EXT; + return COMMIT_FORMATTER.parse(instant + ext); } } + public static Date parseDateFromInstantTime(String instant) throws ParseException { + return parseDateFromInstantTime(instant, true); + } + public static String getInstantForDate(Date instantDate) { return COMMIT_FORMATTER.format(instantDate); } @@ -108,20 +116,27 @@ public static String getInstantForDate(Date instantDate) { /** * Creates an instant string given a valid date-time string. * @param dateString A date-time string in the format yyyy-MM-dd HH:mm:ss[:SSS] + * @param preferMaxDefaultMillis If true and dateString is missing milliseconds then + * append the max number of milliseconds. If false it will be min * @return A timeline instant * @throws ParseException If we cannot parse the date string */ - public static String getInstantForDateString(String dateString) throws ParseException { + public static String getInstantForDateString(String dateString, boolean preferMaxDefaultMillis) throws ParseException { try { return getInstantForDate(MILLIS_GRANULARITY_DATE_FORMATTER.parse(dateString)); } catch (ParseException e) { // Attempt to add the milliseconds in order to complete parsing + String ext = preferMaxDefaultMillis ? MAX_MILLIS_INSTANT_EXT : MIN_MILLIS_INSTANT_EXT; return getInstantForDate(MILLIS_GRANULARITY_DATE_FORMATTER.parse( - String.format("%s:%s", dateString, DEFAULT_MILLIS_EXT) + String.format("%s:%s", dateString, ext) )); } } + public static String getInstantForDateString(String dateString) throws ParseException { + return getInstantForDateString(dateString, true); + } + /** * Returns next instant time in the {@link #COMMIT_FORMATTER} format. * Ensures each instant time is atleast 1 second apart since we create instant times at second granularity diff --git a/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java b/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java index 263988346a3a9..d5acd627192e6 100644 --- a/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java +++ b/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java @@ -81,7 +81,8 @@ void testMedianInstantTime() { String higher = "20210705125921"; String lower = "20210705125806"; String median1 = StreamerUtil.medianInstantTime(higher, lower); - assertThat(median1, is("20210705125843500")); + String expectedMedianWithMillis = "20210705125844499"; + assertThat(median1, is(expectedMedianWithMillis)); // test symmetry assertThrows(IllegalArgumentException.class, () -> StreamerUtil.medianInstantTime(lower, higher), From 150a746006416bc25e73675e794b967ca84e9974 Mon Sep 17 00:00:00 2001 From: Dave Hagman Date: Wed, 27 Oct 2021 17:35:08 -0400 Subject: [PATCH 10/11] Revert "Update test value (which is actually correct). Add ability to override the default number of milliseconds added to second granularity instants" This reverts commit 52b7ea1f9d000003d2a31dd88fc41a35e121827b. --- .../table/timeline/HoodieActiveTimeline.java | 31 +++++-------------- .../apache/hudi/utils/TestStreamerUtil.java | 3 +- 2 files changed, 9 insertions(+), 25 deletions(-) diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java index 6516917d6b1eb..09c8d292f50aa 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/timeline/HoodieActiveTimeline.java @@ -67,11 +67,10 @@ public class HoodieActiveTimeline extends HoodieDefaultTimeline { private static final String MILLIS_GRANULARITY_DATE_FORMAT = "yyyy-MM-dd HH:mm:ss:SSS"; private static final SimpleDateFormat MILLIS_GRANULARITY_DATE_FORMATTER = new SimpleDateFormat(MILLIS_GRANULARITY_DATE_FORMAT); - // Millisecond extensions available when converting from second -> millisecond granularity - // You may prefer the max timestamp as it mimics the current behavior with second granularity - // when performing comparisons such as LESS_THAN_OR_EQUAL_TO but there are times where the min is desirable. - private static final String MIN_MILLIS_INSTANT_EXT = "000"; - private static final String MAX_MILLIS_INSTANT_EXT = "999"; + // The default number of milliseconds that we add if they are not present + // We prefer the max timestamp as it mimics the current behavior with second granularity + // when performing comparisons such as LESS_THAN_OR_EQUAL_TO + private static final String DEFAULT_MILLIS_EXT = "999"; public static final Set VALID_EXTENSIONS_IN_ACTIVE_TIMELINE = new HashSet<>(Arrays.asList( COMMIT_EXTENSION, INFLIGHT_COMMIT_EXTENSION, REQUESTED_COMMIT_EXTENSION, @@ -89,26 +88,19 @@ public class HoodieActiveTimeline extends HoodieDefaultTimeline { /** * Parses the given instant ID to return a date instance. * @param instant The instant ID - * @param preferMaxDefaultMillis If true and dateString is missing milliseconds then - * append the max number of milliseconds. If false it will be min * @return A date * @throws ParseException If the instant ID is malformed */ - public static Date parseDateFromInstantTime(String instant, boolean preferMaxDefaultMillis) throws ParseException { + public static Date parseDateFromInstantTime(String instant) throws ParseException { // Enables backwards compatibility with non-millisecond granularity instants if (isMillisecondGranularity(instant)) { return COMMIT_FORMATTER.parse(instant); } else { // Add milliseconds to the instant in order to parse successfully - String ext = preferMaxDefaultMillis ? MAX_MILLIS_INSTANT_EXT : MIN_MILLIS_INSTANT_EXT; - return COMMIT_FORMATTER.parse(instant + ext); + return COMMIT_FORMATTER.parse(instant + DEFAULT_MILLIS_EXT); } } - public static Date parseDateFromInstantTime(String instant) throws ParseException { - return parseDateFromInstantTime(instant, true); - } - public static String getInstantForDate(Date instantDate) { return COMMIT_FORMATTER.format(instantDate); } @@ -116,27 +108,20 @@ public static String getInstantForDate(Date instantDate) { /** * Creates an instant string given a valid date-time string. * @param dateString A date-time string in the format yyyy-MM-dd HH:mm:ss[:SSS] - * @param preferMaxDefaultMillis If true and dateString is missing milliseconds then - * append the max number of milliseconds. If false it will be min * @return A timeline instant * @throws ParseException If we cannot parse the date string */ - public static String getInstantForDateString(String dateString, boolean preferMaxDefaultMillis) throws ParseException { + public static String getInstantForDateString(String dateString) throws ParseException { try { return getInstantForDate(MILLIS_GRANULARITY_DATE_FORMATTER.parse(dateString)); } catch (ParseException e) { // Attempt to add the milliseconds in order to complete parsing - String ext = preferMaxDefaultMillis ? MAX_MILLIS_INSTANT_EXT : MIN_MILLIS_INSTANT_EXT; return getInstantForDate(MILLIS_GRANULARITY_DATE_FORMATTER.parse( - String.format("%s:%s", dateString, ext) + String.format("%s:%s", dateString, DEFAULT_MILLIS_EXT) )); } } - public static String getInstantForDateString(String dateString) throws ParseException { - return getInstantForDateString(dateString, true); - } - /** * Returns next instant time in the {@link #COMMIT_FORMATTER} format. * Ensures each instant time is atleast 1 second apart since we create instant times at second granularity diff --git a/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java b/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java index d5acd627192e6..263988346a3a9 100644 --- a/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java +++ b/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java @@ -81,8 +81,7 @@ void testMedianInstantTime() { String higher = "20210705125921"; String lower = "20210705125806"; String median1 = StreamerUtil.medianInstantTime(higher, lower); - String expectedMedianWithMillis = "20210705125844499"; - assertThat(median1, is(expectedMedianWithMillis)); + assertThat(median1, is("20210705125843500")); // test symmetry assertThrows(IllegalArgumentException.class, () -> StreamerUtil.medianInstantTime(lower, higher), From 08536c3c278f6da07c5038f59d36f299e0556dc1 Mon Sep 17 00:00:00 2001 From: Dave Hagman Date: Wed, 27 Oct 2021 17:36:39 -0400 Subject: [PATCH 11/11] Update test to use correct value --- .../src/test/java/org/apache/hudi/utils/TestStreamerUtil.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java b/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java index 263988346a3a9..95164f0cd4023 100644 --- a/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java +++ b/hudi-flink/src/test/java/org/apache/hudi/utils/TestStreamerUtil.java @@ -80,8 +80,10 @@ void testInitTableIfNotExists() throws IOException { void testMedianInstantTime() { String higher = "20210705125921"; String lower = "20210705125806"; + String expectedMedianInstant = "20210705125844499"; + String median1 = StreamerUtil.medianInstantTime(higher, lower); - assertThat(median1, is("20210705125843500")); + assertThat(median1, is(expectedMedianInstant)); // test symmetry assertThrows(IllegalArgumentException.class, () -> StreamerUtil.medianInstantTime(lower, higher),