This is an automated email from the ASF dual-hosted git repository.
exceptionfactory pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/main by this push:
new 03134bdcddd NIFI-15989 Improved Purge Threshold time formatting
(#11313)
03134bdcddd is described below
commit 03134bdcdddf0e55f4514175f704405ca83c67ef
Author: David Young <[email protected]>
AuthorDate: Sat Jul 25 13:46:51 2026 -0400
NIFI-15989 Improved Purge Threshold time formatting (#11313)
Signed-off-by: David Handermann <[email protected]>
---
.../java/org/apache/nifi/util/FormatUtils.java | 55 ++++++++++++++++++++++
.../java/org/apache/nifi/util/TestFormatUtils.java | 23 +++++++++
.../nifi/provenance/store/EventStorePartition.java | 4 +-
.../provenance/store/PartitionedEventStore.java | 3 +-
.../provenance/store/WriteAheadStorePartition.java | 16 +++++--
5 files changed, 93 insertions(+), 8 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 5fda79cdc3f..3656473f303 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.processor.DataUnit;
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.ZoneOffset;
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.UnsupportedTemporalTypeException;
+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,54 @@ public class FormatUtils {
final OffsetDateTime offsetDateTime = OffsetDateTime.of(localDateTime,
zoneOffset);
return offsetDateTime.toInstant();
}
+
+ /**
+ * 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(final long value, final
ChronoUnit unit) throws UnsupportedTemporalTypeException {
+ 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(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<String> parts = new ArrayList<>();
+
+ if (days > 0) {
+ parts.add(days + "d");
+ }
+ if (hours > 0 || !parts.isEmpty()) {
+ parts.add(hours + "h");
+ }
+ if (minutes > 0 || !parts.isEmpty()) {
+ parts.add(minutes + "m");
+ }
+ if (seconds > 0 || !parts.isEmpty()) {
+ parts.add(seconds + "s");
+ }
+ 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 b2c1dba37d2..c5e60d8cb77 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
@@ -21,6 +21,7 @@ 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;
@@ -194,4 +195,26 @@ public class TestFormatUtils {
+ 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<Arguments> getDurationValues() {
+ return Stream.of(
+ Arguments.of(Duration.parse("PT0.000000001S"), "1ns"),
+ Arguments.of(Duration.parse("PT0.000000002S"), "2ns"),
+ Arguments.of(Duration.parse("PT1S"), "1s 0ns"),
+ Arguments.of(Duration.parse("PT2S"), "2s 0ns"),
+ Arguments.of(Duration.parse("PT1M"), "1m 0s 0ns"),
+ Arguments.of(Duration.parse("PT2M"), "2m 0s 0ns"),
+ Arguments.of(Duration.parse("PT1H"), "1h 0m 0s 0ns"),
+ Arguments.of(Duration.parse("PT2H"), "2h 0m 0s 0ns"),
+ Arguments.of(Duration.parse("P1D"), "1d 0h 0m 0s 0ns"),
+ Arguments.of(Duration.parse("P35D"), "35d 0h 0m 0s 0ns"),
+ Arguments.of(Duration.parse("P366D"), "366d 0h 0m 0s 0ns")
+ );
+ }
}
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 7e4967ebec7..5e02f9933ac 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
org.apache.nifi.provenance.store.iterator.EventIterator;
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 7d329f5b042..295dc2c1e90 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 org.slf4j.LoggerFactory;
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 @@ public abstract class PartitionedEventStore implements
EventStore {
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 1eec9de83f2..688ba9da577 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,6 +41,8 @@ import java.io.File;
import java.io.FileNotFoundException;
import java.io.IOException;
import java.nio.file.Files;
+import java.time.ZonedDateTime;
+import java.time.temporal.ChronoUnit;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
@@ -490,18 +492,22 @@ public class WriteAheadStorePartition implements
EventStorePartition {
}
@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<File> removed = getEventFilesFromDisk().filter(file ->
file.lastModified() < timeCutoff)
.sorted(DirectoryUtils.SMALLEST_ID_FIRST)
.filter(this::delete)
.collect(Collectors.toList());
+ String thresholdWords = FormatUtils.formatDurationToWords(olderThan,
timeUnit);
+
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);
}
}