From f6c4f26bab1da0921df609596dabc1903d5fa39f Mon Sep 17 00:00:00 2001 From: Archit Goyal Date: Mon, 24 Aug 2026 21:48:42 -0700 Subject: [PATCH 1/5] [FLINK-34246][runtime] Add option to only archive failed jobs to history server Adds jobmanager.archive.only-failed-jobs (default false) to skip archiving finished/canceled/suspended jobs to the history server. --- .../generated/all_jobmanager_section.html | 6 ++++ .../generated/job_manager_configuration.html | 6 ++++ .../configuration/JobManagerOptions.java | 23 +++++++++++++ .../flink/runtime/dispatcher/Dispatcher.java | 14 +++++++- .../DispatcherResourceCleanupTest.java | 34 +++++++++++++++++++ 5 files changed, 82 insertions(+), 1 deletion(-) diff --git a/docs/layouts/shortcodes/generated/all_jobmanager_section.html b/docs/layouts/shortcodes/generated/all_jobmanager_section.html index 5d01a4a4885a5d..a4d2774bca0a32 100644 --- a/docs/layouts/shortcodes/generated/all_jobmanager_section.html +++ b/docs/layouts/shortcodes/generated/all_jobmanager_section.html @@ -86,6 +86,12 @@ String Directory for JobManager to store the archives of completed jobs. + +
jobmanager.archive.only-failed-jobs
+ false + Boolean + Whether to only archive jobs that reached the FAILED terminal state to jobmanager.archive.fs.dir.When enabled, jobs that finished, were canceled, or were suspended are not archived to the history server, reducing the number of files written for large clusters running many short-lived batch jobs. This option has no effect unless jobmanager.archive.fs.dir is configured. +
jobmanager.bind-host
(none) diff --git a/docs/layouts/shortcodes/generated/job_manager_configuration.html b/docs/layouts/shortcodes/generated/job_manager_configuration.html index 75aa1a1d9c821a..3b46b397ced1c2 100644 --- a/docs/layouts/shortcodes/generated/job_manager_configuration.html +++ b/docs/layouts/shortcodes/generated/job_manager_configuration.html @@ -86,6 +86,12 @@ String Directory for JobManager to store the archives of completed jobs. + +
jobmanager.archive.only-failed-jobs
+ false + Boolean + Whether to only archive jobs that reached the FAILED terminal state to jobmanager.archive.fs.dir.When enabled, jobs that finished, were canceled, or were suspended are not archived to the history server, reducing the number of files written for large clusters running many short-lived batch jobs. This option has no effect unless jobmanager.archive.fs.dir is configured. +
jobmanager.bind-host
(none) diff --git a/flink-core/src/main/java/org/apache/flink/configuration/JobManagerOptions.java b/flink-core/src/main/java/org/apache/flink/configuration/JobManagerOptions.java index a185ce22f67773..6e11e83e47404c 100644 --- a/flink-core/src/main/java/org/apache/flink/configuration/JobManagerOptions.java +++ b/flink-core/src/main/java/org/apache/flink/configuration/JobManagerOptions.java @@ -308,6 +308,29 @@ public class JobManagerOptions { .withDescription( "Directory for JobManager to store the archives of completed jobs."); + /** + * Whether only jobs that reached the {@code FAILED} terminal state should be archived to {@link + * #ARCHIVE_DIR}. + */ + @Documentation.Section(Documentation.Sections.ALL_JOB_MANAGER) + public static final ConfigOption ARCHIVE_ON_FAILED_JOBS_ONLY = + key("jobmanager.archive.only-failed-jobs") + .booleanType() + .defaultValue(false) + .withDescription( + Description.builder() + .text( + "Whether to only archive jobs that reached the %s terminal state to %s.", + code("FAILED"), code(ARCHIVE_DIR.key())) + .text( + "When enabled, jobs that finished, were canceled, or were suspended are not " + + "archived to the history server, reducing the number of files written " + + "for large clusters running many short-lived batch jobs. ") + .text( + "This option has no effect unless %s is configured.", + code(ARCHIVE_DIR.key())) + .build()); + /** * @deprecated Use {@link JobManagerOptions#COMPLETED_APPLICATION_STORE_CACHE_SIZE} */ diff --git a/flink-runtime/src/main/java/org/apache/flink/runtime/dispatcher/Dispatcher.java b/flink-runtime/src/main/java/org/apache/flink/runtime/dispatcher/Dispatcher.java index f64645b196bef9..c29be3274f2333 100644 --- a/flink-runtime/src/main/java/org/apache/flink/runtime/dispatcher/Dispatcher.java +++ b/flink-runtime/src/main/java/org/apache/flink/runtime/dispatcher/Dispatcher.java @@ -33,6 +33,7 @@ import org.apache.flink.configuration.Configuration; import org.apache.flink.configuration.DeploymentOptions; import org.apache.flink.configuration.HighAvailabilityOptions; +import org.apache.flink.configuration.JobManagerOptions; import org.apache.flink.configuration.PipelineOptions; import org.apache.flink.configuration.ThreadDumpMode; import org.apache.flink.configuration.WebOptions; @@ -2315,7 +2316,9 @@ protected CompletableFuture jobReachedTerminalState( // do not create an archive for suspended jobs, as this would eventually lead to // multiple archive attempts which we currently do not support CompletableFuture archiveFuture = - archiveExecutionGraphToHistoryServer(executionGraphInfo); + shouldArchiveToHistoryServer(terminalJobStatus) + ? archiveExecutionGraphToHistoryServer(executionGraphInfo) + : CompletableFuture.completedFuture(Acknowledge.get()); return archiveFuture.thenCompose( ignored -> registerGloballyTerminatedJobInJobResultStore(executionGraphInfo)); @@ -2413,6 +2416,15 @@ private void writeToExecutionGraphInfoStore(ExecutionGraphInfo executionGraphInf partialExecutionGraphInfoStore.put(executionGraphInfo.getJobId(), executionGraphInfo); } + /** + * Checks whether a job that reached the given globally terminal state should be archived to the + * history server, honoring {@link JobManagerOptions#ARCHIVE_ON_FAILED_JOBS_ONLY}. + */ + private boolean shouldArchiveToHistoryServer(JobStatus terminalJobStatus) { + return terminalJobStatus == JobStatus.FAILED + || !configuration.get(JobManagerOptions.ARCHIVE_ON_FAILED_JOBS_ONLY); + } + private CompletableFuture archiveExecutionGraphToHistoryServer( ExecutionGraphInfo executionGraphInfo) { diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java index 1d9b3e874754f4..6b65c4977b65cf 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java @@ -22,6 +22,7 @@ import org.apache.flink.api.common.JobID; import org.apache.flink.api.common.JobStatus; import org.apache.flink.configuration.Configuration; +import org.apache.flink.configuration.JobManagerOptions; import org.apache.flink.core.failure.FailureEnricher; import org.apache.flink.core.testutils.OneShotLatch; import org.apache.flink.runtime.application.SingleJobApplication; @@ -736,6 +737,39 @@ public void testNotArchivingSuspendedJobToHistoryServer() throws Exception { assertFalse(isArchived.get()); } + @Test + public void testNotArchivingFinishedJobToHistoryServerWhenOnlyFailedJobsConfigured() + throws Exception { + + final AtomicBoolean isArchived = new AtomicBoolean(false); + + final Configuration configuration = new Configuration(); + configuration.set(JobManagerOptions.ARCHIVE_ON_FAILED_JOBS_ONLY, true); + + final TestingDispatcher.Builder testingDispatcherBuilder = + createTestingDispatcherBuilder() + .setConfiguration(configuration) + .setHistoryServerArchivist( + TestingHistoryServerArchivist.builder() + .setArchiveExecutionGraphFunction( + (executionGraphInfo, applicationId) -> { + isArchived.set(true); + return CompletableFuture.completedFuture( + Acknowledge.get()); + }) + .build()); + + final TestingJobManagerRunnerFactory jobManagerRunnerFactory = + startDispatcherAndSubmitApplication(testingDispatcherBuilder, 0); + + finishJobAndApplication(jobManagerRunnerFactory.takeCreatedJobManagerRunner()); + + assertGlobalCleanupTriggered(jobId); + dispatcher.getJobTerminationFuture(jobId, Duration.ofHours(1)).join(); + + assertFalse(isArchived.get()); + } + private static final class BlockingJobManagerRunnerFactory extends TestingJobMasterServiceLeadershipRunnerFactory { From 42336d671627feb28f4729b642f43752b2594bb4 Mon Sep 17 00:00:00 2001 From: argoyal_Linkedin Date: Thu, 10 Sep 2026 12:02:10 -0700 Subject: [PATCH 2/5] Add missing positive archiving test and fix config option wording --- .../generated/job_manager_configuration.html | 2 +- .../configuration/JobManagerOptions.java | 4 +-- .../DispatcherResourceCleanupTest.java | 35 +++++++++++++++++++ 3 files changed, 38 insertions(+), 3 deletions(-) diff --git a/docs/layouts/shortcodes/generated/job_manager_configuration.html b/docs/layouts/shortcodes/generated/job_manager_configuration.html index 3b46b397ced1c2..1023058a36b599 100644 --- a/docs/layouts/shortcodes/generated/job_manager_configuration.html +++ b/docs/layouts/shortcodes/generated/job_manager_configuration.html @@ -90,7 +90,7 @@
jobmanager.archive.only-failed-jobs
false Boolean - Whether to only archive jobs that reached the FAILED terminal state to jobmanager.archive.fs.dir.When enabled, jobs that finished, were canceled, or were suspended are not archived to the history server, reducing the number of files written for large clusters running many short-lived batch jobs. This option has no effect unless jobmanager.archive.fs.dir is configured. + Whether to only archive jobs that reached the FAILED terminal state to jobmanager.archive.fs.dir. When enabled, jobs that finished or were canceled are not archived to the history server, reducing the number of files written for large clusters running many short-lived batch jobs. This option has no effect unless jobmanager.archive.fs.dir is configured.
jobmanager.bind-host
diff --git a/flink-core/src/main/java/org/apache/flink/configuration/JobManagerOptions.java b/flink-core/src/main/java/org/apache/flink/configuration/JobManagerOptions.java index 6e11e83e47404c..847678a9e45dd4 100644 --- a/flink-core/src/main/java/org/apache/flink/configuration/JobManagerOptions.java +++ b/flink-core/src/main/java/org/apache/flink/configuration/JobManagerOptions.java @@ -320,10 +320,10 @@ public class JobManagerOptions { .withDescription( Description.builder() .text( - "Whether to only archive jobs that reached the %s terminal state to %s.", + "Whether to only archive jobs that reached the %s terminal state to %s. ", code("FAILED"), code(ARCHIVE_DIR.key())) .text( - "When enabled, jobs that finished, were canceled, or were suspended are not " + "When enabled, jobs that finished or were canceled are not " + "archived to the history server, reducing the number of files written " + "for large clusters running many short-lived batch jobs. ") .text( diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java index 6b65c4977b65cf..ea8154334c3193 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java @@ -94,6 +94,7 @@ import static org.hamcrest.Matchers.equalTo; import static org.hamcrest.Matchers.is; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; /** Tests the resource cleanup by the {@link Dispatcher}. */ @@ -770,6 +771,40 @@ public void testNotArchivingFinishedJobToHistoryServerWhenOnlyFailedJobsConfigur assertFalse(isArchived.get()); } + @Test + public void testArchivingFailedJobToHistoryServerWhenOnlyFailedJobsConfigured() + throws Exception { + + final AtomicBoolean isArchived = new AtomicBoolean(false); + + final Configuration configuration = new Configuration(); + configuration.set(JobManagerOptions.ARCHIVE_ON_FAILED_JOBS_ONLY, true); + + final TestingDispatcher.Builder testingDispatcherBuilder = + createTestingDispatcherBuilder() + .setConfiguration(configuration) + .setHistoryServerArchivist( + TestingHistoryServerArchivist.builder() + .setArchiveExecutionGraphFunction( + (executionGraphInfo, applicationId) -> { + isArchived.set(true); + return CompletableFuture.completedFuture( + Acknowledge.get()); + }) + .build()); + + final TestingJobManagerRunnerFactory jobManagerRunnerFactory = + startDispatcherAndSubmitApplication(testingDispatcherBuilder, 0); + + terminateJobWithState( + jobManagerRunnerFactory.takeCreatedJobManagerRunner(), JobStatus.FAILED); + + assertGlobalCleanupTriggered(jobId); + dispatcher.getJobTerminationFuture(jobId, Duration.ofHours(1)).join(); + + assertTrue(isArchived.get()); + } + private static final class BlockingJobManagerRunnerFactory extends TestingJobMasterServiceLeadershipRunnerFactory { From 2edd490772bdb3e4f3cd8a1b7025dc4f7c7197a2 Mon Sep 17 00:00:00 2001 From: argoyal_Linkedin Date: Thu, 10 Sep 2026 15:05:18 -0700 Subject: [PATCH 3/5] Fix new test hanging by marking the application job status as failed --- .../dispatcher/DispatcherResourceCleanupTest.java | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java index ea8154334c3193..7d30289c35b16e 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java @@ -496,6 +496,14 @@ private void suspendJob(TestingJobManagerRunner takeCreatedJobManagerRunner) { terminateJobWithState(takeCreatedJobManagerRunner, JobStatus.SUSPENDED); } + private void failJobAndApplication(TestingJobManagerRunner takeCreatedJobManagerRunner) { + terminateJobWithState(takeCreatedJobManagerRunner, JobStatus.FAILED); + application.jobStatusChanges( + takeCreatedJobManagerRunner.getJobID(), + JobStatus.FAILED, + System.currentTimeMillis()); + } + private void cancelJob(TestingJobManagerRunner takeCreatedJobManagerRunner) { terminateJobWithState(takeCreatedJobManagerRunner, JobStatus.CANCELED); } @@ -796,8 +804,7 @@ public void testArchivingFailedJobToHistoryServerWhenOnlyFailedJobsConfigured() final TestingJobManagerRunnerFactory jobManagerRunnerFactory = startDispatcherAndSubmitApplication(testingDispatcherBuilder, 0); - terminateJobWithState( - jobManagerRunnerFactory.takeCreatedJobManagerRunner(), JobStatus.FAILED); + failJobAndApplication(jobManagerRunnerFactory.takeCreatedJobManagerRunner()); assertGlobalCleanupTriggered(jobId); dispatcher.getJobTerminationFuture(jobId, Duration.ofHours(1)).join(); From 74bb8ac50dfbf9dd4b4238ede3b04211fa57866f Mon Sep 17 00:00:00 2001 From: argoyal2212 Date: Thu, 10 Sep 2026 23:48:05 -0700 Subject: [PATCH 4/5] Set a failure cause on the archived execution graph for failed jobs in the test to avoid a hang --- .../DispatcherResourceCleanupTest.java | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java index 7d30289c35b16e..ee6d146f560150 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/dispatcher/DispatcherResourceCleanupTest.java @@ -33,6 +33,7 @@ import org.apache.flink.runtime.client.JobSubmissionException; import org.apache.flink.runtime.dispatcher.cleanup.TestingResourceCleanerFactory; import org.apache.flink.runtime.executiongraph.ArchivedExecutionGraph; +import org.apache.flink.runtime.executiongraph.ErrorInfo; import org.apache.flink.runtime.executiongraph.JobStatusListener; import org.apache.flink.runtime.heartbeat.HeartbeatServices; import org.apache.flink.runtime.highavailability.HighAvailabilityServices; @@ -510,12 +511,17 @@ private void cancelJob(TestingJobManagerRunner takeCreatedJobManagerRunner) { private void terminateJobWithState( TestingJobManagerRunner takeCreatedJobManagerRunner, JobStatus state) { + final ArchivedExecutionGraphBuilder archivedExecutionGraphBuilder = + new ArchivedExecutionGraphBuilder().setJobID(jobId).setState(state); + + if (state == JobStatus.FAILED) { + archivedExecutionGraphBuilder.setFailureCause( + new ErrorInfo( + new FlinkException("Test job failure"), System.currentTimeMillis())); + } + takeCreatedJobManagerRunner.completeResultFuture( - new ExecutionGraphInfo( - new ArchivedExecutionGraphBuilder() - .setJobID(jobId) - .setState(state) - .build())); + new ExecutionGraphInfo(archivedExecutionGraphBuilder.build())); } private void assertThatNoCleanupWasTriggered() { From b69e6091a4b1a28ce97b423ebb7ba9bcbfd7d0dd Mon Sep 17 00:00:00 2001 From: argoyal2212 Date: Fri, 11 Sep 2026 09:57:51 -0700 Subject: [PATCH 5/5] Regenerate the jobmanager archive only failed jobs config doc so it matches the option description --- docs/layouts/shortcodes/generated/all_jobmanager_section.html | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/layouts/shortcodes/generated/all_jobmanager_section.html b/docs/layouts/shortcodes/generated/all_jobmanager_section.html index a4d2774bca0a32..42eb43eea22dd6 100644 --- a/docs/layouts/shortcodes/generated/all_jobmanager_section.html +++ b/docs/layouts/shortcodes/generated/all_jobmanager_section.html @@ -90,7 +90,7 @@
jobmanager.archive.only-failed-jobs
false Boolean - Whether to only archive jobs that reached the FAILED terminal state to jobmanager.archive.fs.dir.When enabled, jobs that finished, were canceled, or were suspended are not archived to the history server, reducing the number of files written for large clusters running many short-lived batch jobs. This option has no effect unless jobmanager.archive.fs.dir is configured. + Whether to only archive jobs that reached the FAILED terminal state to jobmanager.archive.fs.dir. When enabled, jobs that finished or were canceled are not archived to the history server, reducing the number of files written for large clusters running many short-lived batch jobs. This option has no effect unless jobmanager.archive.fs.dir is configured.
jobmanager.bind-host