From 29c49cd5d2d04bb6a00c4c588c2dc346df0b39e5 Mon Sep 17 00:00:00 2001 From: David Young Date: Tue, 2 Jun 2026 15:16:15 -0400 Subject: [PATCH 1/4] NIFI-15989: Improve purge threshold time formatting Update the purge threshold output to be more human readable. Existing output will show the value in milliseconds which can be difficult to read quickly. Use Apache Commons lang3 to get the threhold in a format like "30 days", "7 days 4 hours", etc. --- .../nifi-persistent-provenance-repository/pom.xml | 4 ++++ .../nifi/provenance/store/WriteAheadStorePartition.java | 7 +++++-- 2 files changed, 9 insertions(+), 2 deletions(-) diff --git a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/pom.xml b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/pom.xml index 7204eca9a595..3693d7263670 100644 --- a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/pom.xml +++ b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/pom.xml @@ -51,5 +51,9 @@ lucene-backward-codecs runtime + + org.apache.commons + commons-lang3 + diff --git a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java index 1eec9de83f2a..10e28ba10f04 100644 --- a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java +++ b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java @@ -17,6 +17,7 @@ package org.apache.nifi.provenance.store; +import org.apache.commons.lang3.time.DurationFormatUtils; import org.apache.nifi.events.EventReporter; import org.apache.nifi.provenance.ProvenanceEventRecord; import org.apache.nifi.provenance.RepositoryConfiguration; @@ -498,10 +499,12 @@ public void purgeOldEvents(final long olderThan, final TimeUnit unit) { .filter(this::delete) .collect(Collectors.toList()); + String thresholdWords = DurationFormatUtils.formatDurationWords(olderThan, true, true); + if (removed.isEmpty()) { - logger.debug("No Provenance Event files that exceed time-based threshold of {} {}", olderThan, unit); + logger.debug("No Provenance Event files that exceed time-based threshold of {}", thresholdWords); } else { - logger.info("Purged {} Provenance Event files from Provenance Repository because the events were older than {} {}: {}", removed.size(), olderThan, unit, removed); + logger.info("Purged {} Provenance Event files from Provenance Repository because the events were older than {} : {}", removed.size(), thresholdWords, removed); } } From 960129d4e942cf694ed05515d3cb3e63e12c8458 Mon Sep 17 00:00:00 2001 From: David Young Date: Mon, 13 Jul 2026 17:56:02 -0400 Subject: [PATCH 2/4] NIFI-15989: Implement a basic duration to words function Update the EventStorePartition interface to use ChronoUnit instead of TimeUnit. TimeUnit is from java.util.concurrent where as ChronoUnit is from java.time.temporal which is overall more appropriate for time based operations. Add tests for duration to words function. --- .../org/apache/nifi/util/FormatUtils.java | 45 +++++++++++++++++++ .../org/apache/nifi/util/TestFormatUtils.java | 25 +++++++++++ .../pom.xml | 4 -- .../provenance/store/EventStorePartition.java | 4 +- .../store/PartitionedEventStore.java | 3 +- .../store/WriteAheadStorePartition.java | 15 ++++--- 6 files changed, 84 insertions(+), 12 deletions(-) diff --git a/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java b/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java index 5fda79cdc3f4..66e218624949 100644 --- a/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java +++ b/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java @@ -20,6 +20,7 @@ import org.apache.nifi.time.DurationFormat; import java.text.NumberFormat; +import java.time.Duration; import java.time.Instant; import java.time.LocalDateTime; import java.time.OffsetDateTime; @@ -28,8 +29,12 @@ import java.time.format.DateTimeFormatter; import java.time.format.DateTimeFormatterBuilder; import java.time.temporal.ChronoField; +import java.time.temporal.ChronoUnit; import java.time.temporal.TemporalAccessor; import java.time.temporal.TemporalQueries; +import java.time.temporal.TemporalUnit; +import java.util.ArrayList; +import java.util.List; import java.util.Locale; import java.util.concurrent.TimeUnit; import java.util.regex.Pattern; @@ -298,4 +303,44 @@ public static Instant parseToInstant(final DateTimeFormatter formatter, final St final OffsetDateTime offsetDateTime = OffsetDateTime.of(localDateTime, zoneOffset); return offsetDateTime.toInstant(); } + + public static String formatDurationToWords(long value, ChronoUnit unit) { + return formatDurationToWords(Duration.ZERO.plus(value, unit)); + } + + /** + * Format a Duration using words (days, hours, minutes, seconds, ns) where all lower units are + * included once a non-zero unit is found. Unit plurality is preserved. + * Maximum resolution is in terms of days. + * + * @param source duration to convert to words + * @return String representation of the given duration + */ + public static String formatDurationToWords(Duration source) { + long days = source.toDaysPart(); + long hours = source.toHoursPart(); + long minutes = source.toMinutesPart(); + int seconds = source.toSecondsPart(); + int nanos = source.toNanosPart(); + + List parts = new ArrayList<>(); + + if (days > 0) { + parts.add(days + (days == 1 ? " day" : " days")); + } + if (hours > 0 || !parts.isEmpty()) { + parts.add(hours + (hours == 1 ? " hour" : " hours")); + } + if (minutes > 0 || !parts.isEmpty()) { + parts.add(minutes + (minutes == 1 ? " minute" : " minutes")); + } + if (seconds > 0 || !parts.isEmpty()) { + parts.add(seconds + (seconds == 1 ? " second" : " seconds")); + } + if(nanos > 0 || !parts.isEmpty()) { + parts.add(nanos + " ns"); + } + + return String.join(" ", parts); + } } diff --git a/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java b/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java index b2c1dba37d2c..46c9f7f9e1d3 100644 --- a/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java +++ b/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java @@ -16,16 +16,19 @@ */ package org.apache.nifi.util; +import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.MethodSource; import java.text.DecimalFormatSymbols; +import java.time.Duration; import java.time.Instant; import java.time.LocalDateTime; import java.time.ZoneId; import java.time.ZoneOffset; import java.time.format.DateTimeFormatter; +import java.time.temporal.ChronoUnit; import java.util.TimeZone; import java.util.concurrent.TimeUnit; import java.util.stream.Stream; @@ -194,4 +197,26 @@ private static Stream getFormatTime() { + TimeUnit.MILLISECONDS.convert(60, TimeUnit.SECONDS) + TimeUnit.MILLISECONDS.convert(1001, TimeUnit.MILLISECONDS), TimeUnit.MILLISECONDS, "1000:01:01.001")); } + + @ParameterizedTest + @MethodSource("getDurationValues") + public void testFormatDurationToWords(Duration duration, String expected) { + assertEquals(expected, FormatUtils.formatDurationToWords(duration)); + } + + private static Stream getDurationValues() { + return Stream.of( + Arguments.of(Duration.parse("PT0.000000001S"), "1 ns"), + Arguments.of(Duration.parse("PT0.000000002S"), "2 ns"), + Arguments.of(Duration.parse("PT1S"), "1 second 0 ns"), + Arguments.of(Duration.parse("PT2S"), "2 seconds 0 ns"), + Arguments.of(Duration.parse("PT1M"), "1 minute 0 seconds 0 ns"), + Arguments.of(Duration.parse("PT2M"), "2 minutes 0 seconds 0 ns"), + Arguments.of(Duration.parse("PT1H"), "1 hour 0 minutes 0 seconds 0 ns"), + Arguments.of(Duration.parse("PT2H"), "2 hours 0 minutes 0 seconds 0 ns"), + Arguments.of(Duration.parse("P1D"), "1 day 0 hours 0 minutes 0 seconds 0 ns"), + Arguments.of(Duration.parse("P35D"), "35 days 0 hours 0 minutes 0 seconds 0 ns"), + Arguments.of(Duration.parse("P366D"), "366 days 0 hours 0 minutes 0 seconds 0 ns") + ); + } } diff --git a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/pom.xml b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/pom.xml index 3693d7263670..7204eca9a595 100644 --- a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/pom.xml +++ b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/pom.xml @@ -51,9 +51,5 @@ lucene-backward-codecs runtime - - org.apache.commons - commons-lang3 - diff --git a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/EventStorePartition.java b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/EventStorePartition.java index 7e4967ebec7f..5e02f9933ac0 100644 --- a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/EventStorePartition.java +++ b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/EventStorePartition.java @@ -23,9 +23,9 @@ import java.io.Closeable; import java.io.IOException; +import java.time.temporal.ChronoUnit; import java.util.List; import java.util.Optional; -import java.util.concurrent.TimeUnit; public interface EventStorePartition extends Closeable { /** @@ -102,7 +102,7 @@ public interface EventStorePartition extends Closeable { * @param olderThan the amount of time for which any event older than this should be removed * @param timeUnit the unit of time that applies to the first argument */ - void purgeOldEvents(long olderThan, TimeUnit timeUnit); + void purgeOldEvents(long olderThan, ChronoUnit timeUnit); /** * Purges some number of events from the partition. The oldest events will be purged. diff --git a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/PartitionedEventStore.java b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/PartitionedEventStore.java index 7d329f5b042c..295dc2c1e900 100644 --- a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/PartitionedEventStore.java +++ b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/PartitionedEventStore.java @@ -32,6 +32,7 @@ import java.io.File; import java.io.IOException; +import java.time.temporal.ChronoUnit; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; @@ -243,7 +244,7 @@ void performMaintenance() { final long maxFileLife = repoConfig.getMaxRecordLife(TimeUnit.MILLISECONDS); for (final EventStorePartition partition : getPartitions()) { try { - partition.purgeOldEvents(maxFileLife, TimeUnit.MILLISECONDS); + partition.purgeOldEvents(maxFileLife, ChronoUnit.MILLIS); } catch (final Exception e) { logger.error("Failed to purge expired events from {}", partition, e); eventReporter.reportEvent(Severity.WARNING, EVENT_CATEGORY, diff --git a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java index 10e28ba10f04..32c0218e30ad 100644 --- a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java +++ b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java @@ -17,7 +17,6 @@ package org.apache.nifi.provenance.store; -import org.apache.commons.lang3.time.DurationFormatUtils; import org.apache.nifi.events.EventReporter; import org.apache.nifi.provenance.ProvenanceEventRecord; import org.apache.nifi.provenance.RepositoryConfiguration; @@ -42,6 +41,10 @@ import java.io.FileNotFoundException; import java.io.IOException; import java.nio.file.Files; +import java.time.Duration; +import java.time.Instant; +import java.time.ZonedDateTime; +import java.time.temporal.ChronoUnit; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -491,15 +494,17 @@ private Optional getPathForEventId(final long id) { } @Override - public void purgeOldEvents(final long olderThan, final TimeUnit unit) { - final long timeCutoff = System.currentTimeMillis() - unit.toMillis(olderThan); - + public void purgeOldEvents(final long olderThan, final ChronoUnit timeUnit) { + // Use ZDT to allow the system to handle a ChronoUnit that is otherwise "estimated" + final long timeCutoff = ZonedDateTime.now() + .minus(olderThan, timeUnit) + .toInstant().toEpochMilli(); final List removed = getEventFilesFromDisk().filter(file -> file.lastModified() < timeCutoff) .sorted(DirectoryUtils.SMALLEST_ID_FIRST) .filter(this::delete) .collect(Collectors.toList()); - String thresholdWords = DurationFormatUtils.formatDurationWords(olderThan, true, true); + String thresholdWords = FormatUtils.formatDurationToWords(olderThan, timeUnit); if (removed.isEmpty()) { logger.debug("No Provenance Event files that exceed time-based threshold of {}", thresholdWords); From 66baa5ba502ac4b97377780d2c965ae08d76c114 Mon Sep 17 00:00:00 2001 From: David Young Date: Tue, 14 Jul 2026 15:12:49 -0400 Subject: [PATCH 3/4] Add javadoc to formatDurationToWords, fix checkstyle Add javadoc to the (long, ChronoUnit) version of the formatDurationToWords function Removed extra imports and fixed formatting --- .../java/org/apache/nifi/util/FormatUtils.java | 16 +++++++++++++--- .../org/apache/nifi/util/TestFormatUtils.java | 2 -- .../store/WriteAheadStorePartition.java | 2 -- 3 files changed, 13 insertions(+), 7 deletions(-) diff --git a/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java b/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java index 66e218624949..cfe68edc26b7 100644 --- a/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java +++ b/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java @@ -32,7 +32,7 @@ import java.time.temporal.ChronoUnit; import java.time.temporal.TemporalAccessor; import java.time.temporal.TemporalQueries; -import java.time.temporal.TemporalUnit; +import java.time.temporal.UnsupportedTemporalTypeException; import java.util.ArrayList; import java.util.List; import java.util.Locale; @@ -304,7 +304,17 @@ public static Instant parseToInstant(final DateTimeFormatter formatter, final St return offsetDateTime.toInstant(); } - public static String formatDurationToWords(long value, ChronoUnit unit) { + /** + * Format a value of ChronoUnit into a word representation. + * Does not handle estimated ChronoUnit values. + * For anything greater than {@link ChronoUnit#DAYS}, use {@link #formatDurationToWords(Duration)}. + * + * @param value the number of the given {@link ChronoUnit} to measure + * @param unit {@link ChronoUnit} that the value represents + * @return String representation of the given value and unit + * @throws UnsupportedTemporalTypeException if the unit is not supported ({@link ChronoUnit#isDurationEstimated()} == true) + */ + public static String formatDurationToWords(long value, ChronoUnit unit) throws UnsupportedTemporalTypeException { return formatDurationToWords(Duration.ZERO.plus(value, unit)); } @@ -337,7 +347,7 @@ public static String formatDurationToWords(Duration source) { if (seconds > 0 || !parts.isEmpty()) { parts.add(seconds + (seconds == 1 ? " second" : " seconds")); } - if(nanos > 0 || !parts.isEmpty()) { + if (nanos > 0 || !parts.isEmpty()) { parts.add(nanos + " ns"); } diff --git a/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java b/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java index 46c9f7f9e1d3..cd8076fad8ac 100644 --- a/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java +++ b/nifi-commons/nifi-utils/src/test/java/org/apache/nifi/util/TestFormatUtils.java @@ -16,7 +16,6 @@ */ package org.apache.nifi.util; -import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.MethodSource; @@ -28,7 +27,6 @@ import java.time.ZoneId; import java.time.ZoneOffset; import java.time.format.DateTimeFormatter; -import java.time.temporal.ChronoUnit; import java.util.TimeZone; import java.util.concurrent.TimeUnit; import java.util.stream.Stream; diff --git a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java index 32c0218e30ad..688ba9da577b 100644 --- a/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java +++ b/nifi-framework-bundle/nifi-framework-extensions/nifi-provenance-repository-bundle/nifi-persistent-provenance-repository/src/main/java/org/apache/nifi/provenance/store/WriteAheadStorePartition.java @@ -41,8 +41,6 @@ import java.io.FileNotFoundException; import java.io.IOException; import java.nio.file.Files; -import java.time.Duration; -import java.time.Instant; import java.time.ZonedDateTime; import java.time.temporal.ChronoUnit; import java.util.ArrayList; From 2bf480b88936ae225aed3e1a285022e1c6918154 Mon Sep 17 00:00:00 2001 From: David Young Date: Mon, 20 Jul 2026 13:13:38 -0400 Subject: [PATCH 4/4] Add final keyword to fields and parameters --- .../java/org/apache/nifi/util/FormatUtils.java | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java b/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java index cfe68edc26b7..6acbe601ee36 100644 --- a/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java +++ b/nifi-commons/nifi-utils/src/main/java/org/apache/nifi/util/FormatUtils.java @@ -314,7 +314,7 @@ public static Instant parseToInstant(final DateTimeFormatter formatter, final St * @return String representation of the given value and unit * @throws UnsupportedTemporalTypeException if the unit is not supported ({@link ChronoUnit#isDurationEstimated()} == true) */ - public static String formatDurationToWords(long value, ChronoUnit unit) throws UnsupportedTemporalTypeException { + public static String formatDurationToWords(final long value, final ChronoUnit unit) throws UnsupportedTemporalTypeException { return formatDurationToWords(Duration.ZERO.plus(value, unit)); } @@ -326,14 +326,14 @@ public static String formatDurationToWords(long value, ChronoUnit unit) throws U * @param source duration to convert to words * @return String representation of the given duration */ - public static String formatDurationToWords(Duration source) { - long days = source.toDaysPart(); - long hours = source.toHoursPart(); - long minutes = source.toMinutesPart(); - int seconds = source.toSecondsPart(); - int nanos = source.toNanosPart(); - - List parts = new ArrayList<>(); + public static String formatDurationToWords(final Duration source) { + final long days = source.toDaysPart(); + final long hours = source.toHoursPart(); + final long minutes = source.toMinutesPart(); + final int seconds = source.toSecondsPart(); + final int nanos = source.toNanosPart(); + + final List parts = new ArrayList<>(); if (days > 0) { parts.add(days + (days == 1 ? " day" : " days"));