committed by
Mahmoud Ben Hassine
parent
708f1c8216
commit
c68da18d3e
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2019 the original author or authors.
|
||||
* Copyright 2013-2023 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.
|
||||
@@ -26,6 +26,7 @@ import java.util.Properties;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.Semaphore;
|
||||
import java.util.stream.Collectors;
|
||||
import javax.batch.operations.BatchRuntimeException;
|
||||
import javax.batch.operations.JobExecutionAlreadyCompleteException;
|
||||
import javax.batch.operations.JobExecutionIsRunningException;
|
||||
@@ -47,6 +48,7 @@ import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.springframework.batch.core.BatchStatus;
|
||||
import org.springframework.batch.core.Entity;
|
||||
import org.springframework.batch.core.ExitStatus;
|
||||
import org.springframework.batch.core.Job;
|
||||
import org.springframework.batch.core.JobParameters;
|
||||
@@ -412,9 +414,11 @@ public class JsrJobOperator implements JobOperator, ApplicationContextAware, Ini
|
||||
List<StepExecution> batchExecutions = new ArrayList<>();
|
||||
|
||||
if(executions != null) {
|
||||
for (org.springframework.batch.core.StepExecution stepExecution : executions) {
|
||||
if(!stepExecution.getStepName().contains(":partition")) {
|
||||
batchExecutions.add(new JsrStepExecution(jobExplorer.getStepExecution(executionId, stepExecution.getId())));
|
||||
Set<Long> stepExecutionIds = executions.stream().map(Entity::getId).collect(Collectors.toSet());
|
||||
org.springframework.batch.core.JobExecution jobExecution = jobExplorer.getJobExecution(executionId);
|
||||
for (org.springframework.batch.core.StepExecution stepExecution : jobExecution.getStepExecutions()) {
|
||||
if(!stepExecution.getStepName().contains(":partition") && stepExecutionIds.contains(stepExecution.getId())) {
|
||||
batchExecutions.add(new JsrStepExecution(stepExecution));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2006-2013 the original author or authors.
|
||||
* Copyright 2006-2023 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.
|
||||
@@ -16,9 +16,12 @@
|
||||
|
||||
package org.springframework.batch.core.partition.support;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
import org.springframework.batch.core.JobExecution;
|
||||
import org.springframework.batch.core.StepExecution;
|
||||
import org.springframework.batch.core.explore.JobExplorer;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
@@ -90,14 +93,16 @@ public class RemoteStepExecutionAggregator implements StepExecutionAggregator, I
|
||||
if (executions == null) {
|
||||
return;
|
||||
}
|
||||
Collection<StepExecution> updates = new ArrayList<>();
|
||||
for (StepExecution stepExecution : executions) {
|
||||
Set<Long> stepExecutionIds = executions.stream().map(stepExecution -> {
|
||||
Long id = stepExecution.getId();
|
||||
Assert.state(id != null, "StepExecution has null id. It must be saved first: " + stepExecution);
|
||||
StepExecution update = jobExplorer.getStepExecution(stepExecution.getJobExecutionId(), id);
|
||||
Assert.state(update != null, "Could not reload StepExecution from JobRepository: " + stepExecution);
|
||||
updates.add(update);
|
||||
}
|
||||
return id;
|
||||
}).collect(Collectors.toSet());
|
||||
JobExecution jobExecution = jobExplorer.getJobExecution(result.getJobExecutionId());
|
||||
Assert.state(jobExecution != null,
|
||||
"Could not load JobExecution from JobRepository for id " + result.getJobExecutionId());
|
||||
List<StepExecution> updates = jobExecution.getStepExecutions().stream()
|
||||
.filter(stepExecution -> stepExecutionIds.contains(stepExecution.getId())).collect(Collectors.toList());
|
||||
delegate.aggregate(result, updates);
|
||||
}
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2018 the original author or authors.
|
||||
* Copyright 2013-2023 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.
|
||||
@@ -403,8 +403,6 @@ public class JsrJobOperatorTests extends AbstractJsrTestCase {
|
||||
jobExecution.addStepExecutions(stepExecutions);
|
||||
|
||||
when(jobExplorer.getJobExecution(5L)).thenReturn(jobExecution);
|
||||
when(jobExplorer.getStepExecution(5L, 1L)).thenReturn(new StepExecution("step1", jobExecution, 1L));
|
||||
when(jobExplorer.getStepExecution(5L, 2L)).thenReturn(new StepExecution("step2", jobExecution, 2L));
|
||||
|
||||
List<javax.batch.runtime.StepExecution> results = jsrJobOperator.getStepExecutions(5L);
|
||||
|
||||
@@ -429,8 +427,6 @@ public class JsrJobOperatorTests extends AbstractJsrTestCase {
|
||||
jobExecution.addStepExecutions(stepExecutions);
|
||||
|
||||
when(jobExplorer.getJobExecution(5L)).thenReturn(jobExecution);
|
||||
when(jobExplorer.getStepExecution(5L, 1L)).thenReturn(new StepExecution("step1", jobExecution, 1L));
|
||||
when(jobExplorer.getStepExecution(5L, 2L)).thenReturn(new StepExecution("step2", jobExecution, 2L));
|
||||
|
||||
List<javax.batch.runtime.StepExecution> results = jsrJobOperator.getStepExecutions(5L);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user