Skip to content
Closed
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 @@ -121,7 +121,7 @@ HAVING COUNT(BJI.JOB_INSTANCE_ID) <= :threshold
List.of(COMPLETED.name(), FAILED.name(), UNKNOWN.name()), "threshold", threshold), Long.class);
}

public Long getNotCompletedPartitionsCount(Long jobExecutionId, String partitionerStepName) {
public Long getNotCompletedPartitionsCount(Long jobExecutionId, String partitionerStepName, String exitCode) {
return namedParameterJdbcTemplate.queryForObject("""
SELECT COUNT(bse.STEP_EXECUTION_ID)
FROM BATCH_STEP_EXECUTION BSE
Expand All @@ -131,7 +131,11 @@ SELECT COUNT(bse.STEP_EXECUTION_ID)
BSE.STEP_NAME <> :stepName
AND
BSE.status <> :status
""", Map.of("jobExecutionId", jobExecutionId, "stepName", partitionerStepName, "status", COMPLETED.name()), Long.class);
AND
BSE.exit_code <> :exitCode
""",
Map.of("jobExecutionId", jobExecutionId, "stepName", partitionerStepName, "status", COMPLETED.name(), "exitCode", exitCode),
Long.class);
}

public void updateJobStatusToFailed(Long stuckJobId, String partitionerStepName) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,12 +20,20 @@

import static org.springframework.transaction.TransactionDefinition.PROPAGATION_REQUIRES_NEW;

import java.time.LocalDateTime;
import java.util.List;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.apache.fineract.infrastructure.core.service.DateUtils;
import org.apache.fineract.infrastructure.jobs.data.partitionedjobs.PartitionedJob;
import org.apache.fineract.infrastructure.jobs.domain.JobExecutionRepository;
import org.springframework.batch.core.BatchStatus;
import org.springframework.batch.core.ExitStatus;
import org.springframework.batch.core.JobExecution;
import org.springframework.batch.core.StepExecution;
import org.springframework.batch.core.explore.JobExplorer;
import org.springframework.batch.core.launch.JobOperator;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.stereotype.Service;
import org.springframework.transaction.TransactionStatus;
import org.springframework.transaction.support.TransactionCallbackWithoutResult;
Expand All @@ -39,14 +47,30 @@ public class StuckJobExecutorServiceImpl implements StuckJobExecutorService {
private final JobExecutionRepository jobExecutionRepository;
private final TransactionTemplate transactionTemplate;
private final JobOperator jobOperator;
private final JobExplorer jobExplorer;
private final JobRepository jobRepository;

private static final String EXIT_CODE_VALUE = "FAILED_AFTER_MANAGER_RESTARTED";

@Override
public void resumeStuckJob(String jobName) {
List<Long> stuckJobIds = getStuckJobIds(jobName);
if (isPartitionedJob(jobName) && areThereStuckJobs(jobName)) {
restartPartitionedJobs(jobName, stuckJobIds);
} else {
restartTaskletJobs(stuckJobIds);
if (areThereStuckJobs(jobName)) {
// Exit Status (Tag) to know If the Job and Steps were marked as Failed
final ExitStatus cleanupTaskExitStatus = new ExitStatus(EXIT_CODE_VALUE);
final List<Long> stuckJobIds = getStuckJobIds(jobName);

// Mark as Failed the Steps Executions (If exists) after Restart
stuckJobIds.forEach(stuckJobId -> verifyBatchStepsExecutions(stuckJobId, cleanupTaskExitStatus));

// Wait the other Job Partitions (If exists) to be completed
if (isPartitionedJob(jobName)) {
restartPartitionedJobs(jobName, stuckJobIds);
} else {
restartTaskletJobs(stuckJobIds);
}

// Mark as Failed the Job Executions (If exists) after Restart
stuckJobIds.forEach(stuckJobId -> verifyBatchJob(stuckJobId, cleanupTaskExitStatus));
}
}

Expand Down Expand Up @@ -108,7 +132,67 @@ private void waitUntilAllPartitionsFinished(Long stuckJobId, String partitionerS
}

private boolean areAllPartitionsCompleted(Long stuckJobId, String partitionerStepName) {
Long notCompletedPartitions = jobExecutionRepository.getNotCompletedPartitionsCount(stuckJobId, partitionerStepName);
Long notCompletedPartitions = jobExecutionRepository.getNotCompletedPartitionsCount(stuckJobId, partitionerStepName,
EXIT_CODE_VALUE);
return notCompletedPartitions == 0L;
}

private void verifyBatchJob(long executionId, ExitStatus exitStatus) {
JobExecution jobExecution = jobExplorer.getJobExecution(executionId);
if (isRunning(jobExecution)) {
log.info("Job Execution {} will be mark as Failed", executionId);
final LocalDateTime now = DateUtils.getLocalDateTimeOfSystem();

// Mark the Job Execution as Failed
jobExecution.setStatus(BatchStatus.FAILED);
jobExecution.setExitStatus(exitStatus);
jobExecution.setEndTime(now);
jobRepository.update(jobExecution);
}
}

private void verifyBatchStepsExecutions(long executionId, ExitStatus exitStatus) {
final JobExecution jobExecution = jobExplorer.getJobExecution(executionId);
if (isRunning(jobExecution)) {
log.info("Step Execution {} will be mark as Failed", executionId);
final LocalDateTime now = DateUtils.getLocalDateTimeOfSystem();

// Mark the Step Executions as Failed
for (StepExecution stepExecution : jobExecution.getStepExecutions()) {
if (isRunning(stepExecution)) {
stepExecution.setStatus(BatchStatus.FAILED);
stepExecution.setExitStatus(exitStatus);
stepExecution.setEndTime(now);
jobRepository.update(stepExecution);
}
}
}
}

private boolean isRunning(JobExecution jobExecution) {
switch (jobExecution.getStatus()) {
case STARTED:
case STARTING:
case STOPPING:
case UNKNOWN:
return true;

default:
return jobExecution.getEndTime() == null;
}
}

private boolean isRunning(StepExecution stepExecution) {
switch (stepExecution.getStatus()) {
case STARTED:
case STARTING:
case STOPPING:
case UNKNOWN:
return true;

default:
return stepExecution.getEndTime() == null;
}
}

}