CodeTrainerMan commented on code in PR #29304:
URL: https://github.com/apache/flink/pull/29304#discussion_r4212021362


##########
flink-runtime/src/main/java/org/apache/flink/runtime/highavailability/nonha/embedded/EmbeddedJobResultStore.java:
##########
@@ -67,7 +105,7 @@ public boolean hasDirtyJobResultEntryInternal(JobID jobId) {
 
     @Override
     public boolean hasCleanJobResultEntryInternal(JobID jobId) {
-        return cleanJobResults.containsKey(jobId);
+        return cleanJobResults.asMap().containsKey(jobId);

Review Comment:
   Good point, thanks. Once an entry is evicted the store no longer knows that 
the job existed, so the same JobID can be accepted again if it arrives under a 
different application - the application-level check still turns the same 
application away because the application result store has no such time limit. 
That matches my understanding of the behaviour, so I made it explicit in the 
description of \job-result-store.clean-job-result.ttl\ (and in the generated 
docs table) to make sure that anyone opting into eviction is aware of it.



##########
flink-runtime/src/main/java/org/apache/flink/runtime/highavailability/nonha/embedded/EmbeddedJobResultStore.java:
##########
@@ -18,28 +18,66 @@
 
 package org.apache.flink.runtime.highavailability.nonha.embedded;
 
+import org.apache.flink.annotation.VisibleForTesting;
 import org.apache.flink.api.common.JobID;
+import org.apache.flink.configuration.Configuration;
 import 
org.apache.flink.runtime.highavailability.AbstractThreadsafeJobResultStore;
 import org.apache.flink.runtime.highavailability.JobResultEntry;
 import org.apache.flink.runtime.highavailability.JobResultStore;
+import org.apache.flink.runtime.highavailability.JobResultStoreOptions;
 import org.apache.flink.runtime.jobmaster.JobResult;
 import org.apache.flink.util.concurrent.Executors;
 
+import org.apache.flink.shaded.guava33.com.google.common.cache.Cache;
+import org.apache.flink.shaded.guava33.com.google.common.cache.CacheBuilder;
+
+import javax.annotation.Nullable;
+
+import java.time.Duration;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.NoSuchElementException;
 import java.util.Set;
+import java.util.concurrent.TimeUnit;
 import java.util.stream.Collectors;
 
+import static org.apache.flink.util.Preconditions.checkNotNull;
+
 /** A thread-safe in-memory implementation of the {@link JobResultStore}. */
 public class EmbeddedJobResultStore extends AbstractThreadsafeJobResultStore {
 
     private final Map<JobID, JobResultEntry> dirtyJobResults = new HashMap<>();
 
-    private final Map<JobID, JobResultEntry> cleanJobResults = new HashMap<>();
+    private final Cache<JobID, JobResultEntry> cleanJobResults;
 
+    /**
+     * Creates a store that retains the clean job results indefinitely. This 
corresponds to the
+     * behaviour this class always had.
+     */
     public EmbeddedJobResultStore() {
+        // a null TTL means "no expiration" which keeps the previous behaviour 
intact
+        this((Duration) null);
+    }
+
+    /**
+     * Creates a store that evicts the clean job results after {@link
+     * JobResultStoreOptions#CLEAN_JOB_RESULT_TTL} has passed. If the option 
is not configured, the
+     * clean job results are retained indefinitely.
+     */
+    public EmbeddedJobResultStore(Configuration configuration) {
+        this(
+                checkNotNull(configuration, "configuration")
+                        .get(JobResultStoreOptions.CLEAN_JOB_RESULT_TTL));
+    }
+
+    @VisibleForTesting
+    EmbeddedJobResultStore(@Nullable Duration cleanJobResultTtl) {
         super(Executors.directExecutor());
+        final CacheBuilder<Object, Object> cacheBuilder = 
CacheBuilder.newBuilder();
+        if (cleanJobResultTtl != null) {
+            cacheBuilder.expireAfterAccess(cleanJobResultTtl.toMillis(), 
TimeUnit.MILLISECONDS);

Review Comment:
   Done. A non-positive TTL is now rejected with an \IllegalArgumentException\ 
instead of building a cache that expires every entry immediately. Covered by 
\EmbeddedJobResultStoreTtlTest#testZeroTtlIsRejected\, 
\#testNegativeTtlIsRejected\ and \#testNonPositiveTtlConfigurationIsRejected\.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to