Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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 {
Expand All @@ -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);
Expand All @@ -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;
}
Expand Down
Original file line number Diff line number Diff line change
@@ -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.
Expand Down Expand Up @@ -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;
Expand All @@ -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
Expand All @@ -81,6 +83,8 @@ public abstract class AbstractJob implements Job, ListableStepLocator, BeanNameA

private ObservationRegistry observationRegistry;

private BatchEventRecorder batchEventRecorder = BatchEventRecorder.DEFAULT;

/**
* Default constructor.
*/
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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());

Expand Down Expand Up @@ -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 + "]";
Expand Down
Original file line number Diff line number Diff line change
@@ -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.
Expand Down Expand Up @@ -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;

Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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) {
Expand All @@ -216,6 +233,8 @@ public static class CommonJobProperties {

private ObservationRegistry observationRegistry;

private BatchEventRecorder batchEventRecorder = BatchEventRecorder.DEFAULT;

private JobParametersIncrementer jobParametersIncrementer;

private JobParametersValidator jobParametersValidator;
Expand All @@ -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;
Expand Down Expand Up @@ -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;
}
Expand Down
Original file line number Diff line number Diff line change
@@ -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.
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -57,6 +57,7 @@
* @author Will Schipp
* @author Mahmoud Ben Hassine
* @author Yejeong Ham
* @author Fabio Molignoni
* @since 6.0
*/
@SuppressWarnings("removal")
Expand All @@ -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();
Expand Down Expand Up @@ -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();
Expand Down
Loading