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]