From 4e7cb6c05fdd6ffd76bd2ab86d9951380dda65f5 Mon Sep 17 00:00:00 2001 From: Fabio Molignoni Date: Sat, 29 Aug 2026 20:49:41 +0200 Subject: [PATCH] GH-5326: Make JFR-based observability opt-in Signed-off-by: Fabio Molignoni --- .../BatchObservabilityBeanPostProcessor.java | 39 ++-- .../batch/core/job/AbstractJob.java | 24 ++- .../core/job/builder/JobBuilderHelper.java | 40 +++- .../support/TaskExecutorJobOperator.java | 19 +- .../observability/BatchEventRecorder.java | 185 ++++++++++++++++++ .../jfr/FlightRecorderBatchEventRecorder.java | 95 +++++++++ .../jfr/events/job/JobExecutionEvent.java | 11 +- .../jfr/events/job/JobLaunchEvent.java | 6 +- .../jfr/events/step/StepExecutionEvent.java | 11 +- .../jfr/events/step/chunk/ChunkScanEvent.java | 11 +- .../step/chunk/ChunkTransactionEvent.java | 11 +- .../events/step/chunk/ChunkWriteEvent.java | 11 +- .../events/step/chunk/ItemProcessEvent.java | 11 +- .../jfr/events/step/chunk/ItemReadEvent.java | 11 +- .../partition/PartitionAggregateEvent.java | 6 +- .../step/partition/PartitionSplitEvent.java | 11 +- .../step/tasklet/TaskletExecutionEvent.java | 11 +- .../batch/core/partition/PartitionStep.java | 14 +- .../batch/core/step/AbstractStep.java | 20 +- .../core/step/builder/JobStepBuilder.java | 3 +- .../core/step/builder/StepBuilderHelper.java | 38 +++- .../core/step/item/ChunkOrientedStep.java | 51 +++-- .../batch/core/step/tasklet/TaskletStep.java | 11 +- ...chObservabilityBeanPostProcessorTests.java | 44 ++++- .../support/TaskExecutorJobOperatorTests.java | 21 ++ .../BatchEventRecorderTests.java | 60 ++++++ ...FlightRecorderBatchEventRecorderTests.java | 127 ++++++++++++ .../pages/spring-batch-observability/jfr.adoc | 14 +- 28 files changed, 826 insertions(+), 90 deletions(-) create mode 100644 spring-batch-core/src/main/java/org/springframework/batch/core/observability/BatchEventRecorder.java create mode 100644 spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/FlightRecorderBatchEventRecorder.java create mode 100644 spring-batch-core/src/test/java/org/springframework/batch/core/observability/BatchEventRecorderTests.java create mode 100644 spring-batch-core/src/test/java/org/springframework/batch/core/observability/jfr/FlightRecorderBatchEventRecorderTests.java diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/configuration/annotation/BatchObservabilityBeanPostProcessor.java b/spring-batch-core/src/main/java/org/springframework/batch/core/configuration/annotation/BatchObservabilityBeanPostProcessor.java index 76a5938e90..e98a2540d4 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/configuration/annotation/BatchObservabilityBeanPostProcessor.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/configuration/annotation/BatchObservabilityBeanPostProcessor.java @@ -23,6 +23,7 @@ import org.springframework.aop.framework.AopProxyUtils; import org.springframework.batch.core.job.AbstractJob; import org.springframework.batch.core.launch.support.TaskExecutorJobOperator; +import org.springframework.batch.core.observability.BatchEventRecorder; import org.springframework.batch.core.step.AbstractStep; import org.springframework.beans.BeansException; import org.springframework.beans.factory.NoSuchBeanDefinitionException; @@ -32,10 +33,11 @@ /** * Bean post processor that configures observable batch artifacts (typically jobs and - * steps) with a Micrometer's observation registry. + * steps) with the available observability components. * * @author Mahmoud Ben Hassine * @author Sanghyuk Jung + * @author Fabio Molignoni * @since 5.0 */ public class BatchObservabilityBeanPostProcessor implements BeanFactoryPostProcessor, BeanPostProcessor { @@ -55,13 +57,13 @@ public Object postProcessAfterInitialization(Object bean, String beanName) throw LOGGER.debug("BeanFactory is not initialized, skipping observation registry injection"); return bean; } - try { - Object target = AopProxyUtils.getSingletonTarget(bean); - if (target == null) { - target = bean; - } - if (target instanceof AbstractJob || target instanceof AbstractStep - || target instanceof TaskExecutorJobOperator) { + Object target = AopProxyUtils.getSingletonTarget(bean); + if (target == null) { + target = bean; + } + if (target instanceof AbstractJob || target instanceof AbstractStep + || target instanceof TaskExecutorJobOperator) { + try { ObservationRegistry observationRegistry = this.beanFactory.getBean(ObservationRegistry.class); if (target instanceof AbstractJob job) { job.setObservationRegistry(observationRegistry); @@ -73,9 +75,24 @@ public Object postProcessAfterInitialization(Object bean, String beanName) throw operator.setObservationRegistry(observationRegistry); } } - } - catch (NoSuchBeanDefinitionException e) { - LOGGER.debug("No Micrometer observation registry found, defaulting to ObservationRegistry.NOOP"); + catch (NoSuchBeanDefinitionException e) { + LOGGER.debug("No Micrometer observation registry found, defaulting to ObservationRegistry.NOOP"); + } + try { + BatchEventRecorder batchEventRecorder = this.beanFactory.getBean(BatchEventRecorder.class); + if (target instanceof AbstractJob job) { + job.setBatchEventRecorder(batchEventRecorder); + } + if (target instanceof AbstractStep step) { + step.setBatchEventRecorder(batchEventRecorder); + } + if (target instanceof TaskExecutorJobOperator operator) { + operator.setBatchEventRecorder(batchEventRecorder); + } + } + catch (NoSuchBeanDefinitionException e) { + LOGGER.debug("No batch event recorder found, defaulting to BatchEventRecorder.DEFAULT"); + } } return bean; } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/job/AbstractJob.java b/spring-batch-core/src/main/java/org/springframework/batch/core/job/AbstractJob.java index e2277013de..5fbb21ea63 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/job/AbstractJob.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/job/AbstractJob.java @@ -1,5 +1,5 @@ /* - * Copyright 2006-2025 the original author or authors. + * Copyright 2006-2026 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -34,8 +34,9 @@ import org.springframework.batch.core.job.parameters.JobParametersValidator; import org.springframework.batch.core.listener.JobExecutionListener; import org.springframework.batch.core.SpringBatchVersion; +import org.springframework.batch.core.observability.BatchEventRecorder; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; import org.springframework.batch.core.observability.BatchMetrics; -import org.springframework.batch.core.observability.jfr.events.job.JobExecutionEvent; import org.springframework.batch.core.observability.micrometer.MicrometerMetrics; import org.springframework.batch.core.step.ListableStepLocator; import org.springframework.batch.core.step.Step; @@ -58,6 +59,7 @@ * @author Lucas Ward * @author Dave Syer * @author Mahmoud Ben Hassine + * @author Fabio Molignoni */ @NullUnmarked // FIXME to remove once default constructors (required by the batch XML // namespace) are removed @@ -81,6 +83,8 @@ public abstract class AbstractJob implements Job, ListableStepLocator, BeanNameA private ObservationRegistry observationRegistry; + private BatchEventRecorder batchEventRecorder = BatchEventRecorder.DEFAULT; + /** * Default constructor. */ @@ -269,8 +273,8 @@ public final void execute(JobExecution execution) throws JobInterruptedException } JobSynchronizationManager.register(execution); - JobExecutionEvent jobExecutionEvent = new JobExecutionEvent(execution.getJobInstance().getJobName(), - execution.getJobInstance().getId(), execution.getId()); + BatchEvent jobExecutionEvent = this.batchEventRecorder.createJobExecutionEvent( + execution.getJobInstance().getJobName(), execution.getJobInstance().getId(), execution.getId()); jobExecutionEvent.begin(); Observation observation = MicrometerMetrics .createObservation(BatchMetrics.METRICS_PREFIX + "job", this.observationRegistry) @@ -336,7 +340,7 @@ public final void execute(JobExecution execution) throws JobInterruptedException execution.setExitStatus(exitStatus.and(newExitStatus)); } stopObservation(execution, observation); - jobExecutionEvent.exitStatus = execution.getExitStatus().getExitCode(); + jobExecutionEvent.setStatus(execution.getExitStatus().getExitCode()); jobExecutionEvent.commit(); execution.setEndTime(LocalDateTime.now()); @@ -424,6 +428,16 @@ public void setObservationRegistry(ObservationRegistry observationRegistry) { this.observationRegistry = observationRegistry; } + /** + * Set the batch event recorder. Defaults to {@link BatchEventRecorder#DEFAULT}. + * @param batchEventRecorder the batch event recorder + * @since 6.1 + */ + public void setBatchEventRecorder(BatchEventRecorder batchEventRecorder) { + Assert.notNull(batchEventRecorder, "BatchEventRecorder must not be null"); + this.batchEventRecorder = batchEventRecorder; + } + @Override public String toString() { return ClassUtils.getShortName(getClass()) + ": [name=" + name + "]"; diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/job/builder/JobBuilderHelper.java b/spring-batch-core/src/main/java/org/springframework/batch/core/job/builder/JobBuilderHelper.java index 54d7859310..1aff21761b 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/job/builder/JobBuilderHelper.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/job/builder/JobBuilderHelper.java @@ -1,5 +1,5 @@ /* - * Copyright 2006-2025 the original author or authors. + * Copyright 2006-2026 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -34,6 +34,7 @@ import org.springframework.batch.core.annotation.BeforeJob; import org.springframework.batch.core.job.AbstractJob; import org.springframework.batch.core.listener.JobListenerFactoryBean; +import org.springframework.batch.core.observability.BatchEventRecorder; import org.springframework.batch.core.repository.JobRepository; import org.springframework.batch.infrastructure.support.ReflectionUtils; @@ -44,6 +45,7 @@ * @author Dave Syer * @author Mahmoud Ben Hassine * @author Taeik Lim + * @author Fabio Molignoni * @since 2.2 */ @NullUnmarked // FIXME to remove once default constructors (required by the batch XML @@ -121,6 +123,20 @@ public B observationRegistry(ObservationRegistry observationRegistry) { return result; } + /** + * Set the batch event recorder for the job. Defaults to + * {@link BatchEventRecorder#DEFAULT}. + * @param batchEventRecorder the batch event recorder + * @return this for fluent chaining + * @since 6.1 + */ + public B batchEventRecorder(BatchEventRecorder batchEventRecorder) { + properties.batchEventRecorder = batchEventRecorder; + @SuppressWarnings("unchecked") + B result = (B) this; + return result; + } + /** * Registers objects using the annotation based listener configuration. * @param listener the object that has a method configured with listener annotation @@ -192,6 +208,7 @@ protected void enhance(AbstractJob job) { if (observationRegistry != null) { job.setObservationRegistry(observationRegistry); } + job.setBatchEventRecorder(properties.getBatchEventRecorder()); Boolean restartable = properties.getRestartable(); if (restartable != null) { @@ -216,6 +233,8 @@ public static class CommonJobProperties { private ObservationRegistry observationRegistry; + private BatchEventRecorder batchEventRecorder = BatchEventRecorder.DEFAULT; + private JobParametersIncrementer jobParametersIncrementer; private JobParametersValidator jobParametersValidator; @@ -228,6 +247,7 @@ public CommonJobProperties(CommonJobProperties properties) { this.restartable = properties.restartable; this.jobRepository = properties.jobRepository; this.observationRegistry = properties.observationRegistry; + this.batchEventRecorder = properties.batchEventRecorder; this.jobExecutionListeners = new LinkedHashSet<>(properties.jobExecutionListeners); this.jobParametersIncrementer = properties.jobParametersIncrementer; this.jobParametersValidator = properties.jobParametersValidator; @@ -265,6 +285,24 @@ public void setObservationRegistry(ObservationRegistry observationRegistry) { this.observationRegistry = observationRegistry; } + /** + * Return the batch event recorder. + * @return the batch event recorder + * @since 6.1 + */ + public BatchEventRecorder getBatchEventRecorder() { + return this.batchEventRecorder; + } + + /** + * Set the batch event recorder. + * @param batchEventRecorder the batch event recorder + * @since 6.1 + */ + public void setBatchEventRecorder(BatchEventRecorder batchEventRecorder) { + this.batchEventRecorder = batchEventRecorder; + } + public String getName() { return name; } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/launch/support/TaskExecutorJobOperator.java b/spring-batch-core/src/main/java/org/springframework/batch/core/launch/support/TaskExecutorJobOperator.java index 1a8df5c97a..f91f44f1d0 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/launch/support/TaskExecutorJobOperator.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/launch/support/TaskExecutorJobOperator.java @@ -1,5 +1,5 @@ /* - * Copyright 2022-2025 the original author or authors. + * Copyright 2022-2026 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -28,7 +28,7 @@ import org.springframework.batch.core.job.parameters.InvalidJobParametersException; import org.springframework.batch.core.launch.JobExecutionNotRunningException; import org.springframework.batch.core.launch.JobOperator; -import org.springframework.batch.core.observability.jfr.events.job.JobLaunchEvent; +import org.springframework.batch.core.observability.BatchEventRecorder; import org.springframework.batch.core.observability.micrometer.MicrometerMetrics; import org.springframework.batch.core.launch.JobExecutionAlreadyRunningException; import org.springframework.batch.core.launch.JobInstanceAlreadyCompleteException; @@ -57,6 +57,7 @@ * @author Will Schipp * @author Mahmoud Ben Hassine * @author Yejeong Ham + * @author Fabio Molignoni * @since 6.0 */ @SuppressWarnings("removal") @@ -66,6 +67,8 @@ public class TaskExecutorJobOperator extends SimpleJobOperator { protected @Nullable ObservationRegistry observationRegistry; + private BatchEventRecorder batchEventRecorder = BatchEventRecorder.DEFAULT; + @Override public void afterPropertiesSet() throws Exception { super.afterPropertiesSet(); @@ -104,12 +107,22 @@ public void setObservationRegistry(ObservationRegistry observationRegistry) { this.observationRegistry = observationRegistry; } + /** + * Set the batch event recorder. Defaults to {@link BatchEventRecorder#DEFAULT}. + * @param batchEventRecorder the batch event recorder + * @since 6.1 + */ + public void setBatchEventRecorder(BatchEventRecorder batchEventRecorder) { + Assert.notNull(batchEventRecorder, "BatchEventRecorder must not be null"); + this.batchEventRecorder = batchEventRecorder; + } + @Override public JobExecution start(Job job, JobParameters jobParameters) throws JobInstanceAlreadyCompleteException, JobExecutionAlreadyRunningException, JobRestartException, InvalidJobParametersException { Assert.notNull(job, "Job must not be null"); Assert.notNull(jobParameters, "JobParameters must not be null"); - new JobLaunchEvent(job.getName(), jobParameters.toString()).commit(); + this.batchEventRecorder.createJobLaunchEvent(job.getName(), jobParameters.toString()).commit(); Observation observation = MicrometerMetrics .createObservation(METRICS_PREFIX + "job.launch.count", this.observationRegistry) .start(); diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/BatchEventRecorder.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/BatchEventRecorder.java new file mode 100644 index 0000000000..e5069c7319 --- /dev/null +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/BatchEventRecorder.java @@ -0,0 +1,185 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.batch.core.observability; + +/** + * Strategy interface for recording Spring Batch execution events. + * + * @author Fabio Molignoni + * @since 6.1 + */ +public interface BatchEventRecorder { + + /** + * Default no-op event recorder. + */ + BatchEventRecorder DEFAULT = new BatchEventRecorder() { + }; + + /** + * Create an event for a job launch request. + * @param jobName the name of the job + * @param jobParameters the job parameters + * @return an event for the job launch request + */ + default BatchEvent createJobLaunchEvent(String jobName, String jobParameters) { + return BatchEvent.DEFAULT; + } + + /** + * Create an event for a job execution. + * @param jobName the name of the job + * @param jobInstanceId the job instance identifier + * @param jobExecutionId the job execution identifier + * @return an event for the job execution + */ + default BatchEvent createJobExecutionEvent(String jobName, long jobInstanceId, long jobExecutionId) { + return BatchEvent.DEFAULT; + } + + /** + * Create an event for a step execution. + * @param stepName the name of the step + * @param jobName the name of the job + * @param stepExecutionId the step execution identifier + * @param jobExecutionId the job execution identifier + * @return an event for the step execution + */ + default BatchEvent createStepExecutionEvent(String stepName, String jobName, long stepExecutionId, + long jobExecutionId) { + return BatchEvent.DEFAULT; + } + + /** + * Create an event for a tasklet execution. + * @param stepName the name of the step + * @param stepExecutionId the step execution identifier + * @param taskletType the fully qualified tasklet type + * @return an event for the tasklet execution + */ + default BatchEvent createTaskletExecutionEvent(String stepName, long stepExecutionId, String taskletType) { + return BatchEvent.DEFAULT; + } + + /** + * Create an event for partition splitting. + * @param stepName the name of the step + * @param stepExecutionId the step execution identifier + * @return an event for partition splitting + */ + default BatchEvent createPartitionSplitEvent(String stepName, long stepExecutionId) { + return BatchEvent.DEFAULT; + } + + /** + * Create an event for partition aggregation. + * @param stepName the name of the step + * @param stepExecutionId the step execution identifier + * @return an event for partition aggregation + */ + default BatchEvent createPartitionAggregateEvent(String stepName, long stepExecutionId) { + return BatchEvent.DEFAULT; + } + + /** + * Create an event for a chunk transaction. + * @param stepName the name of the step + * @param stepExecutionId the step execution identifier + * @return an event for the chunk transaction + */ + default BatchEvent createChunkTransactionEvent(String stepName, long stepExecutionId) { + return BatchEvent.DEFAULT; + } + + /** + * Create an event for a chunk scan. + * @param stepName the name of the step + * @param stepExecutionId the step execution identifier + * @return an event for the chunk scan + */ + default BatchEvent createChunkScanEvent(String stepName, long stepExecutionId) { + return BatchEvent.DEFAULT; + } + + /** + * Create an event for an item read. + * @param stepName the name of the step + * @param stepExecutionId the step execution identifier + * @return an event for the item read + */ + default BatchEvent createItemReadEvent(String stepName, long stepExecutionId) { + return BatchEvent.DEFAULT; + } + + /** + * Create an event for item processing. + * @param stepName the name of the step + * @param stepExecutionId the step execution identifier + * @return an event for item processing + */ + default BatchEvent createItemProcessEvent(String stepName, long stepExecutionId) { + return BatchEvent.DEFAULT; + } + + /** + * Create an event for a chunk write. + * @param stepName the name of the step + * @param stepExecutionId the step execution identifier + * @param itemCount the number of items to write + * @return an event for the chunk write + */ + default BatchEvent createChunkWriteEvent(String stepName, long stepExecutionId, long itemCount) { + return BatchEvent.DEFAULT; + } + + /** + * A recorded batch event. The default implementation is a no-op. + * + * @since 6.1 + */ + interface BatchEvent { + + /** + * Default no-op batch event. + */ + BatchEvent DEFAULT = new BatchEvent() { + }; + + /** Begin the event. */ + default void begin() { + } + + /** + * Set the event outcome. + * @param status the event status + */ + default void setStatus(String status) { + } + + /** + * Set the event count. + * @param count the event count + */ + default void setCount(long count) { + } + + /** Commit the event. */ + default void commit() { + } + + } + +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/FlightRecorderBatchEventRecorder.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/FlightRecorderBatchEventRecorder.java new file mode 100644 index 0000000000..e12cf4f07a --- /dev/null +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/FlightRecorderBatchEventRecorder.java @@ -0,0 +1,95 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.batch.core.observability.jfr; + +import org.springframework.batch.core.observability.BatchEventRecorder; +import org.springframework.batch.core.observability.jfr.events.job.JobExecutionEvent; +import org.springframework.batch.core.observability.jfr.events.job.JobLaunchEvent; +import org.springframework.batch.core.observability.jfr.events.step.StepExecutionEvent; +import org.springframework.batch.core.observability.jfr.events.step.chunk.ChunkScanEvent; +import org.springframework.batch.core.observability.jfr.events.step.chunk.ChunkTransactionEvent; +import org.springframework.batch.core.observability.jfr.events.step.chunk.ChunkWriteEvent; +import org.springframework.batch.core.observability.jfr.events.step.chunk.ItemProcessEvent; +import org.springframework.batch.core.observability.jfr.events.step.chunk.ItemReadEvent; +import org.springframework.batch.core.observability.jfr.events.step.partition.PartitionAggregateEvent; +import org.springframework.batch.core.observability.jfr.events.step.partition.PartitionSplitEvent; +import org.springframework.batch.core.observability.jfr.events.step.tasklet.TaskletExecutionEvent; + +/** + * {@link BatchEventRecorder} implementation for Java Flight Recorder. + * + * @author Fabio Molignoni + * @since 6.1 + */ +public class FlightRecorderBatchEventRecorder implements BatchEventRecorder { + + @Override + public BatchEvent createJobLaunchEvent(String jobName, String jobParameters) { + return new JobLaunchEvent(jobName, jobParameters); + } + + @Override + public BatchEvent createJobExecutionEvent(String jobName, long jobInstanceId, long jobExecutionId) { + return new JobExecutionEvent(jobName, jobInstanceId, jobExecutionId); + } + + @Override + public BatchEvent createStepExecutionEvent(String stepName, String jobName, long stepExecutionId, + long jobExecutionId) { + return new StepExecutionEvent(stepName, jobName, stepExecutionId, jobExecutionId); + } + + @Override + public BatchEvent createTaskletExecutionEvent(String stepName, long stepExecutionId, String taskletType) { + return new TaskletExecutionEvent(stepName, stepExecutionId, taskletType); + } + + @Override + public BatchEvent createPartitionSplitEvent(String stepName, long stepExecutionId) { + return new PartitionSplitEvent(stepName, stepExecutionId); + } + + @Override + public BatchEvent createPartitionAggregateEvent(String stepName, long stepExecutionId) { + return new PartitionAggregateEvent(stepName, stepExecutionId); + } + + @Override + public BatchEvent createChunkTransactionEvent(String stepName, long stepExecutionId) { + return new ChunkTransactionEvent(stepName, stepExecutionId); + } + + @Override + public BatchEvent createChunkScanEvent(String stepName, long stepExecutionId) { + return new ChunkScanEvent(stepName, stepExecutionId); + } + + @Override + public BatchEvent createItemReadEvent(String stepName, long stepExecutionId) { + return new ItemReadEvent(stepName, stepExecutionId); + } + + @Override + public BatchEvent createItemProcessEvent(String stepName, long stepExecutionId) { + return new ItemProcessEvent(stepName, stepExecutionId); + } + + @Override + public BatchEvent createChunkWriteEvent(String stepName, long stepExecutionId, long itemCount) { + return new ChunkWriteEvent(stepName, stepExecutionId, itemCount); + } + +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/job/JobExecutionEvent.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/job/JobExecutionEvent.java index 416c524444..336cb26c5e 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/job/JobExecutionEvent.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/job/JobExecutionEvent.java @@ -20,10 +20,12 @@ import jdk.jfr.Event; import jdk.jfr.Label; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + @Label("Job Execution") @Description("Job Execution Event") @Category({ "Spring Batch", "Job" }) -public class JobExecutionEvent extends Event { +public class JobExecutionEvent extends Event implements BatchEvent { @Label("Job Name") public String jobName; @@ -43,4 +45,9 @@ public JobExecutionEvent(String jobName, long jobInstanceId, long jobExecutionId this.jobExecutionId = jobExecutionId; } -} \ No newline at end of file + @Override + public void setStatus(String status) { + this.exitStatus = status; + } + +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/job/JobLaunchEvent.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/job/JobLaunchEvent.java index 269d099196..d726afa80f 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/job/JobLaunchEvent.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/job/JobLaunchEvent.java @@ -20,10 +20,12 @@ import jdk.jfr.Event; import jdk.jfr.Label; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + @Label("Job Launch Request") @Description("Job Launch Request Event") @Category({ "Spring Batch", "Job" }) -public class JobLaunchEvent extends Event { +public class JobLaunchEvent extends Event implements BatchEvent { @Label("Job Name") public String jobName; @@ -36,4 +38,4 @@ public JobLaunchEvent(String jobName, String jobParameters) { this.jobName = jobName; } -} \ No newline at end of file +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/StepExecutionEvent.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/StepExecutionEvent.java index 6e7784fc79..d849b2e827 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/StepExecutionEvent.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/StepExecutionEvent.java @@ -20,10 +20,12 @@ import jdk.jfr.Event; import jdk.jfr.Label; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + @Label("Step Execution") @Description("Step Execution Event") @Category({ "Spring Batch", "Step" }) -public class StepExecutionEvent extends Event { +public class StepExecutionEvent extends Event implements BatchEvent { @Label("Step Name") public String stepName; @@ -47,4 +49,9 @@ public StepExecutionEvent(String stepName, String jobName, long stepExecutionId, this.jobExecutionId = jobExecutionId; } -} \ No newline at end of file + @Override + public void setStatus(String status) { + this.exitStatus = status; + } + +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkScanEvent.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkScanEvent.java index d5d11e7d1b..57861a79a8 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkScanEvent.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkScanEvent.java @@ -20,10 +20,12 @@ import jdk.jfr.Event; import jdk.jfr.Label; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + @Label("Chunk Scan") @Description("Chunk Scan Event") @Category({ "Spring Batch", "Step", "Chunk" }) -public class ChunkScanEvent extends Event { +public class ChunkScanEvent extends Event implements BatchEvent { @Label("Step Name") public String stepName; @@ -39,4 +41,9 @@ public ChunkScanEvent(String stepName, long stepExecutionId) { this.stepExecutionId = stepExecutionId; } -} \ No newline at end of file + @Override + public void setCount(long count) { + this.skipCount = count; + } + +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkTransactionEvent.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkTransactionEvent.java index 695f3afcfa..190def009b 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkTransactionEvent.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkTransactionEvent.java @@ -20,10 +20,12 @@ import jdk.jfr.Event; import jdk.jfr.Label; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + @Label("Chunk Transaction") @Description("Chunk Transaction Event") @Category({ "Spring Batch", "Step", "Chunk" }) -public class ChunkTransactionEvent extends Event { +public class ChunkTransactionEvent extends Event implements BatchEvent { @Label("Step Name") public String stepName; @@ -39,4 +41,9 @@ public ChunkTransactionEvent(String stepName, long stepExecutionId) { this.stepExecutionId = stepExecutionId; } -} \ No newline at end of file + @Override + public void setStatus(String status) { + this.transactionStatus = status; + } + +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkWriteEvent.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkWriteEvent.java index 6139abb60b..64c08568d6 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkWriteEvent.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ChunkWriteEvent.java @@ -20,10 +20,12 @@ import jdk.jfr.Event; import jdk.jfr.Label; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + @Label("Chunk Write") @Description("Chunk Write Event") @Category({ "Spring Batch", "Step", "Chunk" }) -public class ChunkWriteEvent extends Event { +public class ChunkWriteEvent extends Event implements BatchEvent { @Label("Step Name") public String stepName; @@ -43,4 +45,9 @@ public ChunkWriteEvent(String stepName, long stepExecutionId, long itemCount) { this.stepExecutionId = stepExecutionId; } -} \ No newline at end of file + @Override + public void setStatus(String status) { + this.chunkWriteStatus = status; + } + +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ItemProcessEvent.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ItemProcessEvent.java index 358794dcff..375659c8b4 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ItemProcessEvent.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ItemProcessEvent.java @@ -20,10 +20,12 @@ import jdk.jfr.Event; import jdk.jfr.Label; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + @Label("Item Process") @Description("Item Process Event") @Category({ "Spring Batch", "Step", "Chunk" }) -public class ItemProcessEvent extends Event { +public class ItemProcessEvent extends Event implements BatchEvent { @Label("Step Name") public String stepName; @@ -39,4 +41,9 @@ public ItemProcessEvent(String stepName, long stepExecutionId) { this.stepExecutionId = stepExecutionId; } -} \ No newline at end of file + @Override + public void setStatus(String status) { + this.itemProcessStatus = status; + } + +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ItemReadEvent.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ItemReadEvent.java index 5e55c0de3d..4f6626c5cd 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ItemReadEvent.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/chunk/ItemReadEvent.java @@ -20,10 +20,12 @@ import jdk.jfr.Event; import jdk.jfr.Label; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + @Label("Item Read") @Description("Item Read Event") @Category({ "Spring Batch", "Step", "Chunk" }) -public class ItemReadEvent extends Event { +public class ItemReadEvent extends Event implements BatchEvent { @Label("Step Name") public String stepName; @@ -39,4 +41,9 @@ public ItemReadEvent(String stepName, long stepExecutionId) { this.stepExecutionId = stepExecutionId; } -} \ No newline at end of file + @Override + public void setStatus(String status) { + this.itemReadStatus = status; + } + +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/partition/PartitionAggregateEvent.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/partition/PartitionAggregateEvent.java index f516b08de1..78c1fb8d25 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/partition/PartitionAggregateEvent.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/partition/PartitionAggregateEvent.java @@ -20,10 +20,12 @@ import jdk.jfr.Event; import jdk.jfr.Label; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + @Label("Partition Aggregate") @Description("Partition Aggregate Event") @Category({ "Spring Batch", "Step", "Partition" }) -public class PartitionAggregateEvent extends Event { +public class PartitionAggregateEvent extends Event implements BatchEvent { @Label("Step Name") public String stepName; @@ -36,4 +38,4 @@ public PartitionAggregateEvent(String stepName, long stepExecutionId) { this.stepExecutionId = stepExecutionId; } -} \ No newline at end of file +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/partition/PartitionSplitEvent.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/partition/PartitionSplitEvent.java index 26504edc2f..b1d3d846c3 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/partition/PartitionSplitEvent.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/partition/PartitionSplitEvent.java @@ -20,10 +20,12 @@ import jdk.jfr.Event; import jdk.jfr.Label; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + @Label("Partition Split") @Description("Partition Split Event") @Category({ "Spring Batch", "Step", "Partition" }) -public class PartitionSplitEvent extends Event { +public class PartitionSplitEvent extends Event implements BatchEvent { @Label("Step Name") public String stepName; @@ -39,4 +41,9 @@ public PartitionSplitEvent(String stepName, long stepExecutionId) { this.stepExecutionId = stepExecutionId; } -} \ No newline at end of file + @Override + public void setCount(long count) { + this.partitionCount = count; + } + +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/tasklet/TaskletExecutionEvent.java b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/tasklet/TaskletExecutionEvent.java index 97fbb7b6d7..70f9efcd7c 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/tasklet/TaskletExecutionEvent.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/observability/jfr/events/step/tasklet/TaskletExecutionEvent.java @@ -20,10 +20,12 @@ import jdk.jfr.Event; import jdk.jfr.Label; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + @Label("Tasklet Execution") @Description("Tasklet Execution Event") @Category({ "Spring Batch", "Step", "Tasklet" }) -public class TaskletExecutionEvent extends Event { +public class TaskletExecutionEvent extends Event implements BatchEvent { @Label("Step Name") public String stepName; @@ -43,4 +45,9 @@ public TaskletExecutionEvent(String stepName, long stepExecutionId, String taskl this.stepExecutionId = stepExecutionId; } -} \ No newline at end of file + @Override + public void setStatus(String status) { + this.taskletStatus = status; + } + +} diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/partition/PartitionStep.java b/spring-batch-core/src/main/java/org/springframework/batch/core/partition/PartitionStep.java index 98c754cac9..e4ea860e38 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/partition/PartitionStep.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/partition/PartitionStep.java @@ -1,5 +1,5 @@ /* - * Copyright 2006-2025 the original author or authors. + * Copyright 2006-2026 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -18,8 +18,7 @@ import org.springframework.batch.core.BatchStatus; import org.springframework.batch.core.job.JobExecutionException; -import org.springframework.batch.core.observability.jfr.events.step.partition.PartitionAggregateEvent; -import org.springframework.batch.core.observability.jfr.events.step.partition.PartitionSplitEvent; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; import org.springframework.batch.core.repository.JobRepository; import org.springframework.batch.core.step.Step; import org.springframework.batch.core.step.StepExecution; @@ -38,6 +37,7 @@ * * @author Dave Syer * @author Mahmoud Ben Hassine + * @author Fabio Molignoni * @since 2.0 */ @NullUnmarked // FIXME to remove once default constructors (required by the batch XML @@ -114,17 +114,17 @@ protected void doExecute(StepExecution stepExecution) throws Exception { stepExecution.getExecutionContext().put(STEP_TYPE_KEY, this.getClass().getName()); // Split execution into partitions and wait for task completion - PartitionSplitEvent partitionSplitEvent = new PartitionSplitEvent(stepExecution.getStepName(), + BatchEvent partitionSplitEvent = this.batchEventRecorder.createPartitionSplitEvent(stepExecution.getStepName(), stepExecution.getId()); partitionSplitEvent.begin(); Collection executions = partitionHandler.handle(stepExecutionSplitter, stepExecution); - partitionSplitEvent.partitionCount = executions.size(); + partitionSplitEvent.setCount(executions.size()); stepExecution.upgradeStatus(BatchStatus.COMPLETED); partitionSplitEvent.commit(); // aggregate the results of the executions - PartitionAggregateEvent partitionAggregateEvent = new PartitionAggregateEvent(stepExecution.getStepName(), - stepExecution.getId()); + BatchEvent partitionAggregateEvent = this.batchEventRecorder + .createPartitionAggregateEvent(stepExecution.getStepName(), stepExecution.getId()); partitionAggregateEvent.begin(); stepExecutionAggregator.aggregate(stepExecution, executions); partitionAggregateEvent.commit(); diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java index d7979364f1..2dc8f5d3da 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java @@ -36,8 +36,9 @@ import org.springframework.batch.core.launch.NoSuchJobException; import org.springframework.batch.core.launch.support.ExitCodeMapper; import org.springframework.batch.core.listener.CompositeStepExecutionListener; +import org.springframework.batch.core.observability.BatchEventRecorder; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; import org.springframework.batch.core.observability.BatchMetrics; -import org.springframework.batch.core.observability.jfr.events.step.StepExecutionEvent; import org.springframework.batch.core.observability.micrometer.MicrometerMetrics; import org.springframework.batch.core.repository.JobRepository; import org.springframework.batch.core.scope.context.StepSynchronizationManager; @@ -59,6 +60,7 @@ * @author Chris Schaefer * @author Mahmoud Ben Hassine * @author Jinwoo Bae + * @author Fabio Molignoni */ // FIXME remove once default constructors (required by the XML namespace) are removed @NullUnmarked @@ -78,6 +80,8 @@ public abstract class AbstractStep implements StoppableStep, InitializingBean, B protected ObservationRegistry observationRegistry; + protected BatchEventRecorder batchEventRecorder = BatchEventRecorder.DEFAULT; + /** * Create a new {@link AbstractStep}. * @deprecated since 6.0 for removal in 7.0. Use {@link #AbstractStep(JobRepository)} @@ -217,7 +221,7 @@ public final void execute(StepExecution stepExecution) if (logger.isDebugEnabled()) { logger.debug("Executing: id=" + stepExecution.getId()); } - StepExecutionEvent stepExecutionEvent = new StepExecutionEvent(stepExecution.getStepName(), + BatchEvent stepExecutionEvent = this.batchEventRecorder.createStepExecutionEvent(stepExecution.getStepName(), stepExecution.getJobExecution().getJobInstance().getJobName(), stepExecution.getId(), stepExecution.getJobExecutionId()); stepExecutionEvent.begin(); @@ -320,7 +324,7 @@ public final void execute(StepExecution stepExecution) stepExecution.getJobExecution().getJobInstance().getJobName()), e); } - stepExecutionEvent.exitStatus = stepExecution.getExitStatus().getExitCode(); + stepExecutionEvent.setStatus(stepExecution.getExitStatus().getExitCode()); stepExecutionEvent.commit(); stopObservation(stepExecution, observation); stepExecution.setExitStatus(exitStatus); @@ -468,4 +472,14 @@ public void setObservationRegistry(ObservationRegistry observationRegistry) { this.observationRegistry = observationRegistry; } + /** + * Set the batch event recorder. Defaults to {@link BatchEventRecorder#DEFAULT}. + * @param batchEventRecorder the batch event recorder + * @since 6.1 + */ + public void setBatchEventRecorder(BatchEventRecorder batchEventRecorder) { + Assert.notNull(batchEventRecorder, "BatchEventRecorder must not be null"); + this.batchEventRecorder = batchEventRecorder; + } + } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/builder/JobStepBuilder.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/builder/JobStepBuilder.java index 586be304b4..c96fcaf262 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/builder/JobStepBuilder.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/builder/JobStepBuilder.java @@ -1,5 +1,5 @@ /* - * Copyright 2006-2025 the original author or authors. + * Copyright 2006-2026 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -102,6 +102,7 @@ public Step build() { TaskExecutorJobOperator jobOperator = new TaskExecutorJobOperator(); jobOperator.setJobRepository(getJobRepository()); jobOperator.setJobRegistry(new MapJobRegistry()); + jobOperator.setBatchEventRecorder(this.properties.getBatchEventRecorder()); try { jobOperator.afterPropertiesSet(); } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/builder/StepBuilderHelper.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/builder/StepBuilderHelper.java index da00172ffe..dac8ff14a0 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/builder/StepBuilderHelper.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/builder/StepBuilderHelper.java @@ -1,5 +1,5 @@ /* - * Copyright 2006-2025 the original author or authors. + * Copyright 2006-2026 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -30,6 +30,7 @@ import org.springframework.batch.core.annotation.AfterStep; import org.springframework.batch.core.annotation.BeforeStep; import org.springframework.batch.core.listener.StepListenerFactoryBean; +import org.springframework.batch.core.observability.BatchEventRecorder; import org.springframework.batch.core.repository.JobRepository; import org.springframework.batch.core.step.AbstractStep; import org.springframework.batch.infrastructure.support.ReflectionUtils; @@ -42,6 +43,7 @@ * @author Michael Minella * @author Taeik Lim * @author Mahmoud Ben Hassine + * @author Fabio Molignoni * @since 2.2 */ // FIXME remove once default constructors (required by the XML namespace) are removed @@ -93,6 +95,18 @@ public B startLimit(int startLimit) { return self(); } + /** + * Set the batch event recorder for the step. Defaults to + * {@link BatchEventRecorder#DEFAULT}. + * @param batchEventRecorder the batch event recorder + * @return this for fluent chaining + * @since 6.1 + */ + public B batchEventRecorder(BatchEventRecorder batchEventRecorder) { + properties.batchEventRecorder = batchEventRecorder; + return self(); + } + /** * Registers objects using the annotation based listener configuration. * @param listener the object that has a method configured with listener annotation @@ -143,6 +157,7 @@ protected void enhance(AbstractStep step) { if (observationRegistry != null) { step.setObservationRegistry(observationRegistry); } + step.setBatchEventRecorder(properties.getBatchEventRecorder()); Boolean allowStartIfComplete = properties.allowStartIfComplete; if (allowStartIfComplete != null) { @@ -171,6 +186,8 @@ public static class CommonStepProperties { private ObservationRegistry observationRegistry = ObservationRegistry.NOOP; + private BatchEventRecorder batchEventRecorder = BatchEventRecorder.DEFAULT; + public CommonStepProperties() { } @@ -180,6 +197,7 @@ public CommonStepProperties(CommonStepProperties properties) { this.allowStartIfComplete = properties.allowStartIfComplete; this.jobRepository = properties.jobRepository; this.observationRegistry = properties.observationRegistry; + this.batchEventRecorder = properties.batchEventRecorder; this.stepExecutionListeners = new ArrayList<>(properties.stepExecutionListeners); } @@ -199,6 +217,24 @@ public void setObservationRegistry(ObservationRegistry observationRegistry) { this.observationRegistry = observationRegistry; } + /** + * Return the batch event recorder. + * @return the batch event recorder + * @since 6.1 + */ + public BatchEventRecorder getBatchEventRecorder() { + return this.batchEventRecorder; + } + + /** + * Set the batch event recorder. + * @param batchEventRecorder the batch event recorder + * @since 6.1 + */ + public void setBatchEventRecorder(BatchEventRecorder batchEventRecorder) { + this.batchEventRecorder = batchEventRecorder; + } + public String getName() { return name; } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedStep.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedStep.java index c85fd19d9f..ea36dca134 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedStep.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/item/ChunkOrientedStep.java @@ -36,12 +36,8 @@ import org.springframework.batch.core.listener.ItemReadListener; import org.springframework.batch.core.listener.ItemWriteListener; import org.springframework.batch.core.listener.SkipListener; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; import org.springframework.batch.core.observability.BatchMetrics; -import org.springframework.batch.core.observability.jfr.events.step.chunk.ChunkScanEvent; -import org.springframework.batch.core.observability.jfr.events.step.chunk.ChunkTransactionEvent; -import org.springframework.batch.core.observability.jfr.events.step.chunk.ChunkWriteEvent; -import org.springframework.batch.core.observability.jfr.events.step.chunk.ItemProcessEvent; -import org.springframework.batch.core.observability.jfr.events.step.chunk.ItemReadEvent; import org.springframework.batch.core.scope.context.StepContext; import org.springframework.batch.core.scope.context.StepSynchronizationManager; import org.springframework.batch.core.step.StepContribution; @@ -93,6 +89,7 @@ * @author Minchul Son * @author Yanming Zhou * @author Taeik Lim + * @author Fabio Molignoni * @since 6.0 */ public class ChunkOrientedStep extends AbstractStep { @@ -391,8 +388,8 @@ protected void doExecute(StepExecution stepExecution) throws Exception { while (this.chunkTracker.get().moreItems() && !interrupted(stepExecution)) { // process next chunk in its own transaction this.transactionTemplate.executeWithoutResult(transactionStatus -> { - ChunkTransactionEvent chunkTransactionEvent = new ChunkTransactionEvent(stepExecution.getStepName(), - stepExecution.getId()); + BatchEvent chunkTransactionEvent = this.batchEventRecorder + .createChunkTransactionEvent(stepExecution.getStepName(), stepExecution.getId()); chunkTransactionEvent.begin(); StepContribution contribution = stepExecution.createStepContribution(); processNextChunk(transactionStatus, contribution, stepExecution); @@ -404,7 +401,7 @@ protected void doExecute(StepExecution stepExecution) throws Exception { // (eg JpaTransactionManager) has marked it as globally rollback-only // (eg after a JPA flush failure) but not locally rollback-only. transactionStatus.setRollbackOnly(); - chunkTransactionEvent.transactionStatus = BatchMetrics.STATUS_ROLLED_BACK; + chunkTransactionEvent.setStatus(BatchMetrics.STATUS_ROLLED_BACK); chunkTransactionEvent.commit(); return; } @@ -412,7 +409,7 @@ protected void doExecute(StepExecution stepExecution) throws Exception { this.compositeItemStream.update(stepExecution.getExecutionContext()); getJobRepository().updateExecutionContext(stepExecution); getJobRepository().update(stepExecution); - chunkTransactionEvent.transactionStatus = BatchMetrics.STATUS_COMMITTED; + chunkTransactionEvent.setStatus(BatchMetrics.STATUS_COMMITTED); chunkTransactionEvent.commit(); }); } @@ -442,13 +439,13 @@ private void processChunkConcurrently(TransactionStatus status, StepContribution logger.info("Executing scan in new transaction after rollback"); ScanItem scanItem = tracker.pollNextScanItem(); if (scanItem != null) { - ChunkScanEvent chunkScanEvent = new ChunkScanEvent(stepExecution.getStepName(), - stepExecution.getId()); + BatchEvent chunkScanEvent = this.batchEventRecorder + .createChunkScanEvent(stepExecution.getStepName(), stepExecution.getId()); chunkScanEvent.begin(); // no chunk listener callbacks here: ChunkListener is not called in // concurrent steps scan(scanItem, contribution, status); - chunkScanEvent.skipCount = contribution.getSkipCount(); + chunkScanEvent.setCount(contribution.getSkipCount()); chunkScanEvent.commit(); } if (!tracker.hasPendingScanItems()) { @@ -538,15 +535,15 @@ private void processChunkSequentially(TransactionStatus status, StepContribution logger.info("Executing scan in new transaction after rollback"); ScanItem scanItem = tracker.pollNextScanItem(); if (scanItem != null) { - ChunkScanEvent chunkScanEvent = new ChunkScanEvent(stepExecution.getStepName(), - stepExecution.getId()); + BatchEvent chunkScanEvent = this.batchEventRecorder + .createChunkScanEvent(stepExecution.getStepName(), stepExecution.getId()); chunkScanEvent.begin(); compositeChunkListener.beforeChunk(new Chunk<>(scanItem.input())); Chunk singleItemChunk = scan(scanItem, contribution, status); if (!status.isRollbackOnly()) { compositeChunkListener.afterChunk(singleItemChunk); } - chunkScanEvent.skipCount = contribution.getSkipCount(); + chunkScanEvent.setCount(contribution.getSkipCount()); chunkScanEvent.commit(); } if (!tracker.hasPendingScanItems()) { @@ -643,8 +640,8 @@ private Chunk readChunk(StepContribution contribution) throws Exception { } private @Nullable I readItem(StepContribution contribution) throws Exception { - ItemReadEvent itemReadEvent = new ItemReadEvent(contribution.getStepExecution().getStepName(), - contribution.getStepExecution().getId()); + BatchEvent itemReadEvent = this.batchEventRecorder.createItemReadEvent( + contribution.getStepExecution().getStepName(), contribution.getStepExecution().getId()); String fullyQualifiedMetricName = BatchMetrics.METRICS_PREFIX + "item.read"; Observation observation = Observation.createNotStarted(fullyQualifiedMetricName, this.observationRegistry) .lowCardinalityKeyValue(fullyQualifiedMetricName + ".job.name", @@ -664,7 +661,7 @@ private Chunk readChunk(StepContribution contribution) throws Exception { contribution.incrementReadCount(); this.compositeItemReadListener.afterRead(item); } - itemReadEvent.itemReadStatus = BatchMetrics.STATUS_SUCCESS; + itemReadEvent.setStatus(BatchMetrics.STATUS_SUCCESS); observation.lowCardinalityKeyValue(fullyQualifiedMetricName + ".status", BatchMetrics.STATUS_SUCCESS); } catch (Exception exception) { @@ -675,7 +672,7 @@ private Chunk readChunk(StepContribution contribution) throws Exception { else { throw exception; } - itemReadEvent.itemReadStatus = BatchMetrics.STATUS_FAILURE; + itemReadEvent.setStatus(BatchMetrics.STATUS_FAILURE); observation.lowCardinalityKeyValue(fullyQualifiedMetricName + ".status", BatchMetrics.STATUS_FAILURE); observation.error(exception); } @@ -734,8 +731,8 @@ private Chunk processChunk(Chunk chunk, StepContribution contribution, Lis } private @Nullable O processItem(I item, StepContribution contribution) throws Exception { - ItemProcessEvent itemProcessEvent = new ItemProcessEvent(contribution.getStepExecution().getStepName(), - contribution.getStepExecution().getId()); + BatchEvent itemProcessEvent = this.batchEventRecorder.createItemProcessEvent( + contribution.getStepExecution().getStepName(), contribution.getStepExecution().getId()); String fullyQualifiedMetricName = METRICS_PREFIX + "item.process"; Observation observation = Observation.createNotStarted(fullyQualifiedMetricName, this.observationRegistry) .lowCardinalityKeyValue(fullyQualifiedMetricName + ".job.name", @@ -752,7 +749,7 @@ private Chunk processChunk(Chunk chunk, StepContribution contribution, Lis contribution.incrementFilterCount(); } this.compositeItemProcessListener.afterProcess(item, processedItem); - itemProcessEvent.itemProcessStatus = BatchMetrics.STATUS_SUCCESS; + itemProcessEvent.setStatus(BatchMetrics.STATUS_SUCCESS); observation.lowCardinalityKeyValue(fullyQualifiedMetricName + ".status", BatchMetrics.STATUS_SUCCESS); } catch (Exception exception) { @@ -763,7 +760,7 @@ private Chunk processChunk(Chunk chunk, StepContribution contribution, Lis else { throw exception; } - itemProcessEvent.itemProcessStatus = BatchMetrics.STATUS_FAILURE; + itemProcessEvent.setStatus(BatchMetrics.STATUS_FAILURE); observation.lowCardinalityKeyValue(fullyQualifiedMetricName + ".status", BatchMetrics.STATUS_FAILURE); observation.error(exception); } @@ -820,8 +817,8 @@ private void doSkipInProcess(I item, RetryException retryException, StepContribu private void writeChunk(Chunk chunk, List> scanItems, StepContribution contribution, TransactionStatus status) throws Exception { - ChunkWriteEvent chunkWriteEvent = new ChunkWriteEvent(contribution.getStepExecution().getStepName(), - contribution.getStepExecution().getId(), chunk.size()); + BatchEvent chunkWriteEvent = this.batchEventRecorder.createChunkWriteEvent( + contribution.getStepExecution().getStepName(), contribution.getStepExecution().getId(), chunk.size()); String fullyQualifiedMetricName = METRICS_PREFIX + "chunk.write"; Observation observation = Observation.createNotStarted(fullyQualifiedMetricName, this.observationRegistry) .lowCardinalityKeyValue(fullyQualifiedMetricName + ".job.name", @@ -844,12 +841,12 @@ private void writeChunk(Chunk chunk, List> scanItems, StepCont } contribution.incrementWriteCount(chunk.size()); this.compositeItemWriteListener.afterWrite(chunk); - chunkWriteEvent.chunkWriteStatus = BatchMetrics.STATUS_SUCCESS; + chunkWriteEvent.setStatus(BatchMetrics.STATUS_SUCCESS); observation.lowCardinalityKeyValue(fullyQualifiedMetricName + ".status", BatchMetrics.STATUS_SUCCESS); } catch (Exception exception) { this.compositeItemWriteListener.onWriteError(exception, chunk); - chunkWriteEvent.chunkWriteStatus = BatchMetrics.STATUS_FAILURE; + chunkWriteEvent.setStatus(BatchMetrics.STATUS_FAILURE); observation.lowCardinalityKeyValue(fullyQualifiedMetricName + ".status", BatchMetrics.STATUS_FAILURE); observation.error(exception); diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/tasklet/TaskletStep.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/tasklet/TaskletStep.java index 4f9c8cc395..cab6fda88b 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/tasklet/TaskletStep.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/tasklet/TaskletStep.java @@ -1,5 +1,5 @@ /* - * Copyright 2006-2025 the original author or authors. + * Copyright 2006-2026 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -22,7 +22,7 @@ import org.springframework.batch.core.BatchStatus; import org.springframework.batch.core.listener.ChunkListener; import org.springframework.batch.core.job.JobInterruptedException; -import org.springframework.batch.core.observability.jfr.events.step.tasklet.TaskletExecutionEvent; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; import org.springframework.batch.core.step.StepContribution; import org.springframework.batch.core.step.StepExecution; import org.springframework.batch.core.listener.StepExecutionListener; @@ -74,6 +74,7 @@ * @author Michael Minella * @author Will Schipp * @author Mahmoud Ben Hassine + * @author Fabio Molignoni */ // FIXME remove once default constructors (required by the XML namespace) are removed @NullUnmarked @@ -249,8 +250,8 @@ protected void doExecute(StepExecution stepExecution) throws Exception { String taskletType = tasklet.getClass().getName(); stepExecution.getExecutionContext().put(TASKLET_TYPE_KEY, taskletType); stepExecution.getExecutionContext().put(STEP_TYPE_KEY, this.getClass().getName()); - TaskletExecutionEvent taskletExecutionEvent = new TaskletExecutionEvent(stepExecution.getStepName(), - stepExecution.getId(), taskletType); + BatchEvent taskletExecutionEvent = this.batchEventRecorder + .createTaskletExecutionEvent(stepExecution.getStepName(), stepExecution.getId(), taskletType); taskletExecutionEvent.begin(); stream.update(stepExecution.getExecutionContext()); getJobRepository().updateExecutionContext(stepExecution); @@ -293,7 +294,7 @@ public RepeatStatus doInChunkContext(RepeatContext repeatContext, ChunkContext c }); - taskletExecutionEvent.taskletStatus = stepExecution.getExitStatus().getExitCode(); + taskletExecutionEvent.setStatus(stepExecution.getExitStatus().getExitCode()); taskletExecutionEvent.commit(); } diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/configuration/annotation/BatchObservabilityBeanPostProcessorTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/configuration/annotation/BatchObservabilityBeanPostProcessorTests.java index 7ac9520228..3e23ae8ccf 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/configuration/annotation/BatchObservabilityBeanPostProcessorTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/configuration/annotation/BatchObservabilityBeanPostProcessorTests.java @@ -19,17 +19,23 @@ import org.junit.jupiter.api.Test; import org.springframework.aop.framework.ProxyFactory; +import org.springframework.batch.core.job.SimpleJob; import org.springframework.batch.core.launch.JobOperator; import org.springframework.batch.core.launch.support.TaskExecutorJobOperator; +import org.springframework.batch.core.observability.BatchEventRecorder; +import org.springframework.batch.core.repository.support.ResourcelessJobRepository; +import org.springframework.batch.core.step.tasklet.TaskletStep; import org.springframework.beans.factory.support.DefaultListableBeanFactory; import org.springframework.test.util.ReflectionTestUtils; import static org.junit.jupiter.api.Assertions.assertSame; +import static org.mockito.Mockito.mock; /** * Test class for {@link BatchObservabilityBeanPostProcessor}. * * @author Sanghyuk Jung + * @author Fabio Molignoni */ class BatchObservabilityBeanPostProcessorTests { @@ -37,12 +43,43 @@ class BatchObservabilityBeanPostProcessorTests { private final ObservationRegistry observationRegistry = ObservationRegistry.create(); + private final BatchEventRecorder batchEventRecorder = mock(); + private final BatchObservabilityBeanPostProcessor postProcessor = new BatchObservabilityBeanPostProcessor(); @Test - void observationRegistryShouldBeSetOnJobOperator() { + void observabilityComponentsShouldBeSetOnJob() { + this.beanFactory.registerSingleton("observationRegistry", this.observationRegistry); + this.beanFactory.registerSingleton("batchEventRecorder", this.batchEventRecorder); + this.postProcessor.postProcessBeanFactory(this.beanFactory); + SimpleJob job = new SimpleJob("job"); + job.setJobRepository(new ResourcelessJobRepository()); + + this.postProcessor.postProcessAfterInitialization(job, "job"); + + assertSame(this.observationRegistry, ReflectionTestUtils.getField(job, "observationRegistry")); + assertSame(this.batchEventRecorder, ReflectionTestUtils.getField(job, "batchEventRecorder")); + } + + @Test + void observabilityComponentsShouldBeSetOnStep() { + this.beanFactory.registerSingleton("observationRegistry", this.observationRegistry); + this.beanFactory.registerSingleton("batchEventRecorder", this.batchEventRecorder); + this.postProcessor.postProcessBeanFactory(this.beanFactory); + TaskletStep step = new TaskletStep(new ResourcelessJobRepository()); + step.setName("step"); + + this.postProcessor.postProcessAfterInitialization(step, "step"); + + assertSame(this.observationRegistry, ReflectionTestUtils.getField(step, "observationRegistry")); + assertSame(this.batchEventRecorder, ReflectionTestUtils.getField(step, "batchEventRecorder")); + } + + @Test + void observabilityComponentsShouldBeSetOnJobOperator() { // given this.beanFactory.registerSingleton("observationRegistry", this.observationRegistry); + this.beanFactory.registerSingleton("batchEventRecorder", this.batchEventRecorder); this.postProcessor.postProcessBeanFactory(this.beanFactory); TaskExecutorJobOperator jobOperator = new TaskExecutorJobOperator(); @@ -51,12 +88,14 @@ void observationRegistryShouldBeSetOnJobOperator() { // then assertSame(this.observationRegistry, ReflectionTestUtils.getField(jobOperator, "observationRegistry")); + assertSame(this.batchEventRecorder, ReflectionTestUtils.getField(jobOperator, "batchEventRecorder")); } @Test - void observationRegistryShouldBeSetOnProxiedJobOperator() { + void observabilityComponentsShouldBeSetOnProxiedJobOperator() { // given this.beanFactory.registerSingleton("observationRegistry", this.observationRegistry); + this.beanFactory.registerSingleton("batchEventRecorder", this.batchEventRecorder); this.postProcessor.postProcessBeanFactory(this.beanFactory); TaskExecutorJobOperator target = new TaskExecutorJobOperator(); ProxyFactory proxyFactory = new ProxyFactory(); @@ -70,6 +109,7 @@ void observationRegistryShouldBeSetOnProxiedJobOperator() { // then assertSame(this.observationRegistry, ReflectionTestUtils.getField(target, "observationRegistry")); + assertSame(this.batchEventRecorder, ReflectionTestUtils.getField(target, "batchEventRecorder")); } } diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/launch/support/TaskExecutorJobOperatorTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/launch/support/TaskExecutorJobOperatorTests.java index 8e50831951..521c2cb6e4 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/launch/support/TaskExecutorJobOperatorTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/launch/support/TaskExecutorJobOperatorTests.java @@ -32,6 +32,8 @@ import org.springframework.batch.core.launch.NoSuchJobException; import org.springframework.batch.core.launch.JobExecutionAlreadyRunningException; import org.springframework.batch.core.launch.JobInstanceAlreadyCompleteException; +import org.springframework.batch.core.observability.BatchEventRecorder; +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; import org.springframework.batch.core.repository.JobRepository; import org.springframework.batch.core.launch.JobRestartException; import org.springframework.batch.core.repository.support.JdbcJobRepositoryFactoryBean; @@ -45,6 +47,11 @@ import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseType; import org.springframework.jdbc.support.JdbcTransactionManager; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + /** * @author Dave Syer * @author Will Schipp @@ -99,6 +106,20 @@ void testStart() throws JobInstanceAlreadyCompleteException, NoSuchJobException, Assertions.assertEquals(BatchStatus.COMPLETED, jobExecution.getStatus()); } + @Test + void testStartRecordsJobLaunchEvent() throws Exception { + BatchEventRecorder eventRecorder = mock(); + BatchEvent jobLaunchEvent = mock(); + when(eventRecorder.createJobLaunchEvent(anyString(), anyString())).thenReturn(jobLaunchEvent); + jobOperator.setBatchEventRecorder(eventRecorder); + JobParameters jobParameters = new JobParameters(); + + jobOperator.start(job, jobParameters); + + verify(eventRecorder).createJobLaunchEvent("job", jobParameters.toString()); + verify(jobLaunchEvent).commit(); + } + @Test void testRestart() throws Exception { Tasklet tasklet = new Tasklet() { diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/observability/BatchEventRecorderTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/observability/BatchEventRecorderTests.java new file mode 100644 index 0000000000..e702a2941a --- /dev/null +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/observability/BatchEventRecorderTests.java @@ -0,0 +1,60 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.batch.core.observability; + +import org.junit.jupiter.api.Test; + +import org.springframework.batch.core.observability.BatchEventRecorder.BatchEvent; + +import static org.junit.jupiter.api.Assertions.assertAll; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertSame; + +/** + * Tests for {@link BatchEventRecorder}. + * + * @author Fabio Molignoni + */ +class BatchEventRecorderTests { + + @Test + void defaultRecorderShouldCreateNoOpEvents() { + BatchEventRecorder recorder = BatchEventRecorder.DEFAULT; + + assertAll(() -> assertSame(BatchEvent.DEFAULT, recorder.createJobLaunchEvent("job", "parameters")), + () -> assertSame(BatchEvent.DEFAULT, recorder.createJobExecutionEvent("job", 1, 2)), + () -> assertSame(BatchEvent.DEFAULT, recorder.createStepExecutionEvent("step", "job", 3, 2)), + () -> assertSame(BatchEvent.DEFAULT, recorder.createTaskletExecutionEvent("step", 3, "tasklet")), + () -> assertSame(BatchEvent.DEFAULT, recorder.createPartitionSplitEvent("step", 3)), + () -> assertSame(BatchEvent.DEFAULT, recorder.createPartitionAggregateEvent("step", 3)), + () -> assertSame(BatchEvent.DEFAULT, recorder.createChunkTransactionEvent("step", 3)), + () -> assertSame(BatchEvent.DEFAULT, recorder.createChunkScanEvent("step", 3)), + () -> assertSame(BatchEvent.DEFAULT, recorder.createItemReadEvent("step", 3)), + () -> assertSame(BatchEvent.DEFAULT, recorder.createItemProcessEvent("step", 3)), + () -> assertSame(BatchEvent.DEFAULT, recorder.createChunkWriteEvent("step", 3, 4))); + } + + @Test + void defaultEventShouldAcceptLifecycleCallbacks() { + assertDoesNotThrow(() -> { + BatchEvent.DEFAULT.begin(); + BatchEvent.DEFAULT.setStatus("COMPLETED"); + BatchEvent.DEFAULT.setCount(1); + BatchEvent.DEFAULT.commit(); + }); + } + +} diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/observability/jfr/FlightRecorderBatchEventRecorderTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/observability/jfr/FlightRecorderBatchEventRecorderTests.java new file mode 100644 index 0000000000..6bdc940c8a --- /dev/null +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/observability/jfr/FlightRecorderBatchEventRecorderTests.java @@ -0,0 +1,127 @@ +/* + * Copyright 2026 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.batch.core.observability.jfr; + +import org.junit.jupiter.api.Test; + +import org.springframework.batch.core.observability.jfr.events.job.JobExecutionEvent; +import org.springframework.batch.core.observability.jfr.events.job.JobLaunchEvent; +import org.springframework.batch.core.observability.jfr.events.step.StepExecutionEvent; +import org.springframework.batch.core.observability.jfr.events.step.chunk.ChunkScanEvent; +import org.springframework.batch.core.observability.jfr.events.step.chunk.ChunkTransactionEvent; +import org.springframework.batch.core.observability.jfr.events.step.chunk.ChunkWriteEvent; +import org.springframework.batch.core.observability.jfr.events.step.chunk.ItemProcessEvent; +import org.springframework.batch.core.observability.jfr.events.step.chunk.ItemReadEvent; +import org.springframework.batch.core.observability.jfr.events.step.partition.PartitionAggregateEvent; +import org.springframework.batch.core.observability.jfr.events.step.partition.PartitionSplitEvent; +import org.springframework.batch.core.observability.jfr.events.step.tasklet.TaskletExecutionEvent; + +import static org.junit.jupiter.api.Assertions.assertAll; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; + +/** + * Tests for {@link FlightRecorderBatchEventRecorder}. + * + * @author Fabio Molignoni + */ +class FlightRecorderBatchEventRecorderTests { + + private final FlightRecorderBatchEventRecorder recorder = new FlightRecorderBatchEventRecorder(); + + @Test + void shouldCreateJobEvents() { + JobLaunchEvent jobLaunchEvent = assertInstanceOf(JobLaunchEvent.class, + this.recorder.createJobLaunchEvent("job", "parameters")); + JobExecutionEvent jobExecutionEvent = assertInstanceOf(JobExecutionEvent.class, + this.recorder.createJobExecutionEvent("job", 1, 2)); + jobExecutionEvent.setStatus("COMPLETED"); + + assertAll(() -> assertEquals("job", jobLaunchEvent.jobName), + () -> assertEquals("parameters", jobLaunchEvent.jobParameters), + () -> assertEquals("job", jobExecutionEvent.jobName), + () -> assertEquals(1, jobExecutionEvent.jobInstanceId), + () -> assertEquals(2, jobExecutionEvent.jobExecutionId), + () -> assertEquals("COMPLETED", jobExecutionEvent.exitStatus)); + } + + @Test + void shouldCreateStepAndTaskletEvents() { + StepExecutionEvent stepExecutionEvent = assertInstanceOf(StepExecutionEvent.class, + this.recorder.createStepExecutionEvent("step", "job", 3, 2)); + TaskletExecutionEvent taskletExecutionEvent = assertInstanceOf(TaskletExecutionEvent.class, + this.recorder.createTaskletExecutionEvent("step", 3, "tasklet")); + stepExecutionEvent.setStatus("COMPLETED"); + taskletExecutionEvent.setStatus("COMPLETED"); + + assertAll(() -> assertEquals("step", stepExecutionEvent.stepName), + () -> assertEquals("job", stepExecutionEvent.jobName), + () -> assertEquals(3, stepExecutionEvent.stepExecutionId), + () -> assertEquals(2, stepExecutionEvent.jobExecutionId), + () -> assertEquals("COMPLETED", stepExecutionEvent.exitStatus), + () -> assertEquals("step", taskletExecutionEvent.stepName), + () -> assertEquals(3, taskletExecutionEvent.stepExecutionId), + () -> assertEquals("tasklet", taskletExecutionEvent.taskletType), + () -> assertEquals("COMPLETED", taskletExecutionEvent.taskletStatus)); + } + + @Test + void shouldCreatePartitionEvents() { + PartitionSplitEvent partitionSplitEvent = assertInstanceOf(PartitionSplitEvent.class, + this.recorder.createPartitionSplitEvent("step", 3)); + PartitionAggregateEvent partitionAggregateEvent = assertInstanceOf(PartitionAggregateEvent.class, + this.recorder.createPartitionAggregateEvent("step", 3)); + partitionSplitEvent.setCount(4); + + assertAll(() -> assertEquals("step", partitionSplitEvent.stepName), + () -> assertEquals(3, partitionSplitEvent.stepExecutionId), + () -> assertEquals(4, partitionSplitEvent.partitionCount), + () -> assertEquals("step", partitionAggregateEvent.stepName), + () -> assertEquals(3, partitionAggregateEvent.stepExecutionId)); + } + + @Test + void shouldCreateChunkEvents() { + ChunkTransactionEvent transactionEvent = assertInstanceOf(ChunkTransactionEvent.class, + this.recorder.createChunkTransactionEvent("step", 3)); + ChunkScanEvent scanEvent = assertInstanceOf(ChunkScanEvent.class, + this.recorder.createChunkScanEvent("step", 3)); + ItemReadEvent readEvent = assertInstanceOf(ItemReadEvent.class, this.recorder.createItemReadEvent("step", 3)); + ItemProcessEvent processEvent = assertInstanceOf(ItemProcessEvent.class, + this.recorder.createItemProcessEvent("step", 3)); + ChunkWriteEvent writeEvent = assertInstanceOf(ChunkWriteEvent.class, + this.recorder.createChunkWriteEvent("step", 3, 4)); + transactionEvent.setStatus("COMMITTED"); + scanEvent.setCount(1); + readEvent.setStatus("SUCCESS"); + processEvent.setStatus("SUCCESS"); + writeEvent.setStatus("SUCCESS"); + + assertAll(() -> assertEquals("step", transactionEvent.stepName), + () -> assertEquals(3, transactionEvent.stepExecutionId), + () -> assertEquals("COMMITTED", transactionEvent.transactionStatus), + () -> assertEquals("step", scanEvent.stepName), () -> assertEquals(3, scanEvent.stepExecutionId), + () -> assertEquals(1, scanEvent.skipCount), () -> assertEquals("step", readEvent.stepName), + () -> assertEquals(3, readEvent.stepExecutionId), + () -> assertEquals("SUCCESS", readEvent.itemReadStatus), + () -> assertEquals("step", processEvent.stepName), () -> assertEquals(3, processEvent.stepExecutionId), + () -> assertEquals("SUCCESS", processEvent.itemProcessStatus), + () -> assertEquals("step", writeEvent.stepName), () -> assertEquals(3, writeEvent.stepExecutionId), + () -> assertEquals(4, writeEvent.itemCount), + () -> assertEquals("SUCCESS", writeEvent.chunkWriteStatus)); + } + +} diff --git a/spring-batch-docs/modules/ROOT/pages/spring-batch-observability/jfr.adoc b/spring-batch-docs/modules/ROOT/pages/spring-batch-observability/jfr.adoc index 411f51b6d9..cb32776171 100644 --- a/spring-batch-docs/modules/ROOT/pages/spring-batch-observability/jfr.adoc +++ b/spring-batch-docs/modules/ROOT/pages/spring-batch-observability/jfr.adoc @@ -3,11 +3,21 @@ As of version 6, Spring Batch provides support for Java Flight Recorder (JFR) to help you monitor and troubleshoot batch jobs. JFR is a low-overhead, event-based profiling tool built into the Java Virtual Machine (JVM) that allows developers to collect detailed information about the performance and behavior of their applications. -JFR can be enabled by adding the following JVM options when starting your Spring Batch application: +JFR can be enabled by registering a `FlightRecorderBatchEventRecorder` in the application context: + +[source, java] +---- +@Bean +public BatchEventRecorder batchEventRecorder() { + return new FlightRecorderBatchEventRecorder(); +} +---- + +Start the application with recording enabled to capture the events: [source, bash] ---- java -XX:StartFlightRecording:filename=my-batch-job.jfr,dumponexit=true -jar my-batch-job.jar ---- -Once JFR is enabled, Spring Batch will automatically create JFR events for key batch processing activities, such as job and step executions, item reads and writes, as well as transaction boundaries. These events can be viewed and analyzed using tools such as Java Mission Control (JMC) or other JFR-compatible tools. \ No newline at end of file +Once JFR is enabled, Spring Batch will automatically create JFR events for key batch processing activities, such as job and step executions, item reads and writes, as well as transaction boundaries. These events can be viewed and analyzed using tools such as Java Mission Control (JMC) or other JFR-compatible tools.