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


##########
flink-runtime/src/main/java/org/apache/flink/runtime/highavailability/nonha/AbstractNonHaServices.java:
##########
@@ -58,7 +59,15 @@ public abstract class AbstractNonHaServices implements 
HighAvailabilityServices
     private boolean shutdown;
 
     public AbstractNonHaServices() {
-        this.jobResultStore = new EmbeddedJobResultStore();
+        this(new EmbeddedJobResultStore());

Review Comment:
   Should the setting reach the store used by MiniCluster?
   
   The change passes the configuration through StandaloneHaServices. However, I 
think EmbeddedHaServices may still be using the original constructor, which 
doesn't receive the configuration. That class is used when high availability is 
turned off (HighAvailabilityServicesUtils) and in MiniCluster.
   
   In both places the configuration is available but doesn't seem to be passed 
on. 



##########
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:
   worth validating that the value is positive?



##########
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:
   Could a finished job be run a second time after its entry is removed?
   
   IIUC, the job result store does two jobs. It records results, and it stops 
the same job from being run twice. When a job is submitted, the Dispatcher asks 
the store whether it already has an entry for that job ID. If it does, the job 
is turned away (Dispatcher#isInGloballyTerminalState, which calls 
hasJobResultEntryAsync).
   
   With this change, once the time limit has passed, the clean entry is removed 
and the store no longer knows the job existed. If the same job ID is submitted 
again after that, it looks like the job would be accepted and run again.
   
   There is a separate check at the application level, and the application 
result store has no time limit, so submitting the same application again is 
still turned away. The case I have in mind is the same job ID arriving under a 
different application. If that's known/intentional, maybe it's just worth a 
sentence in the option's description to make it clear?



-- 
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