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]