| Line 16... |
Line 16... |
| 16 |
|
16 |
|
| 17 |
import java.time.LocalDateTime;
|
17 |
import java.time.LocalDateTime;
|
| 18 |
import java.time.format.DateTimeFormatter;
|
18 |
import java.time.format.DateTimeFormatter;
|
| 19 |
import java.util.List;
|
19 |
import java.util.List;
|
| 20 |
import java.util.Map;
|
20 |
import java.util.Map;
|
| - |
|
21 |
import java.util.stream.Collectors;
|
| 21 |
|
22 |
|
| 22 |
@Service
|
23 |
@Service
|
| 23 |
public class CronBatchService {
|
24 |
public class CronBatchService {
|
| 24 |
|
25 |
|
| 25 |
private static final Logger LOGGER = LogManager.getLogger(CronBatchService.class);
|
26 |
private static final Logger LOGGER = LogManager.getLogger(CronBatchService.class);
|
| 26 |
private static final String[] TECHNOLOGY_EMAIL = {"amit.gupta@smartdukaan.com"};
|
27 |
private static final String[] TECHNOLOGY_EMAIL = {"amit.gupta@smartdukaan.com"};
|
| 27 |
private static final DateTimeFormatter DTF = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
|
28 |
private static final DateTimeFormatter DTF = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
|
| - |
|
29 |
/** Cap on how many never-attempted items the alert lists by name. */
|
| - |
|
30 |
private static final int PENDING_LIST_CAP = 50;
|
| 28 |
|
31 |
|
| 29 |
@Autowired
|
32 |
@Autowired
|
| 30 |
private CronBatchRepository cronBatchRepository;
|
33 |
private CronBatchRepository cronBatchRepository;
|
| 31 |
|
34 |
|
| 32 |
@Autowired
|
35 |
@Autowired
|
| Line 60... |
Line 63... |
| 60 |
|
63 |
|
| 61 |
LOGGER.info("Created batch {} for job {} with {} items", batch.getId(), jobName, fofoIdPartnerNameMap.size());
|
64 |
LOGGER.info("Created batch {} for job {} with {} items", batch.getId(), jobName, fofoIdPartnerNameMap.size());
|
| 62 |
return batch;
|
65 |
return batch;
|
| 63 |
}
|
66 |
}
|
| 64 |
|
67 |
|
| - |
|
68 |
/**
|
| - |
|
69 |
* REQUIRED — deliberately NOT REQUIRES_NEW.
|
| - |
|
70 |
*
|
| - |
|
71 |
* Every caller is the last statement of a REQUIRES_NEW per-item work method, so this
|
| - |
|
72 |
* joins that work transaction and the SUCCESS marker commits (or rolls back) with the
|
| - |
|
73 |
* work it describes. Under REQUIRES_NEW the marker committed on its own, so work that
|
| - |
|
74 |
* later failed at flush/commit — a constraint violation, a wallet deadlock — was left
|
| - |
|
75 |
* recorded as a successful payout.
|
| - |
|
76 |
*/
|
| 65 |
@Transactional(propagation = Propagation.REQUIRES_NEW)
|
77 |
@Transactional(propagation = Propagation.REQUIRED)
|
| 66 |
public void markItemSuccess(int batchId, int fofoId) {
|
78 |
public void markItemSuccess(int batchId, int fofoId) {
|
| 67 |
List<CronBatchItem> items = cronBatchItemRepository.selectByBatchIdAndStatus(batchId, CronBatchItemStatus.PENDING);
|
79 |
CronBatchItem item = cronBatchItemRepository.selectByBatchIdAndFofoId(batchId, fofoId);
|
| 68 |
for (CronBatchItem item : items) {
|
- |
|
| 69 |
if (item.getFofoId() == fofoId) {
|
80 |
if (item != null && item.getStatus() == CronBatchItemStatus.PENDING) {
|
| 70 |
item.markSuccess();
|
81 |
item.markSuccess();
|
| 71 |
break;
|
- |
|
| 72 |
}
|
- |
|
| 73 |
}
|
82 |
}
|
| 74 |
}
|
83 |
}
|
| 75 |
|
84 |
|
| - |
|
85 |
/**
|
| - |
|
86 |
* REQUIRES_NEW — the opposite of markItemSuccess, and intentionally so: this is called
|
| - |
|
87 |
* from the orchestrator's catch block after the work transaction has already rolled
|
| - |
|
88 |
* back, so the failure record must commit independently of it.
|
| - |
|
89 |
*/
|
| 76 |
@Transactional(propagation = Propagation.REQUIRES_NEW)
|
90 |
@Transactional(propagation = Propagation.REQUIRES_NEW)
|
| 77 |
public void markItemFailed(int batchId, int fofoId, String errorMessage) {
|
91 |
public void markItemFailed(int batchId, int fofoId, String errorMessage) {
|
| 78 |
List<CronBatchItem> items = cronBatchItemRepository.selectByBatchIdAndStatus(batchId, CronBatchItemStatus.PENDING);
|
92 |
CronBatchItem item = cronBatchItemRepository.selectByBatchIdAndFofoId(batchId, fofoId);
|
| 79 |
for (CronBatchItem item : items) {
|
- |
|
| 80 |
if (item.getFofoId() == fofoId) {
|
93 |
if (item != null && item.getStatus() == CronBatchItemStatus.PENDING) {
|
| 81 |
item.markFailed(errorMessage);
|
94 |
item.markFailed(errorMessage);
|
| 82 |
break;
|
- |
|
| 83 |
}
|
- |
|
| 84 |
}
|
95 |
}
|
| 85 |
}
|
96 |
}
|
| 86 |
|
97 |
|
| - |
|
98 |
/**
|
| - |
|
99 |
* RUNNING batches older than {@code staleHours} — a run that died before finalizing.
|
| - |
|
100 |
* Returns ids only; the caller must finalize each through the proxy (see
|
| - |
|
101 |
* BatchScheduledTasks.reapStaleBatches) so REQUIRES_NEW actually applies.
|
| - |
|
102 |
*/
|
| - |
|
103 |
@Transactional(readOnly = true)
|
| - |
|
104 |
public List<Integer> selectStaleRunningBatchIds(int staleHours) {
|
| - |
|
105 |
return cronBatchRepository.selectStaleRunning(LocalDateTime.now().minusHours(staleHours))
|
| - |
|
106 |
.stream().map(CronBatch::getId).collect(Collectors.toList());
|
| - |
|
107 |
}
|
| - |
|
108 |
|
| - |
|
109 |
/**
|
| - |
|
110 |
* Closes a batch using the item rows as the source of truth.
|
| - |
|
111 |
*
|
| - |
|
112 |
* successCount is COUNTED, never derived as (total - failures). Deriving it reported
|
| - |
|
113 |
* every item that never reported at all — the JVM killed mid-loop, a lost marker write —
|
| - |
|
114 |
* as a successful payout, and suppressed the alert, because failureCount stayed 0.
|
| - |
|
115 |
*
|
| - |
|
116 |
* Items still PENDING were never attempted. That leaves the batch INCOMPLETE (not
|
| - |
|
117 |
* PARTIAL_FAILURE, which means "attempted and failed") and always raises an email.
|
| - |
|
118 |
* PENDING rows are left as they are: "never tried" and "tried and failed" are
|
| - |
|
119 |
* different facts and the retry decision depends on which one it was.
|
| - |
|
120 |
*/
|
| 87 |
@Transactional(propagation = Propagation.REQUIRES_NEW)
|
121 |
@Transactional(propagation = Propagation.REQUIRES_NEW)
|
| 88 |
public void finalizeBatch(int batchId) {
|
122 |
public BatchOutcome finalizeBatch(int batchId) {
|
| 89 |
CronBatch batch = cronBatchRepository.selectById(batchId);
|
123 |
CronBatch batch = cronBatchRepository.selectById(batchId);
|
| - |
|
124 |
if (batch == null) {
|
| - |
|
125 |
LOGGER.warn("finalizeBatch called for unknown batch {}", batchId);
|
| - |
|
126 |
return null;
|
| - |
|
127 |
}
|
| 90 |
List<CronBatchItem> failedItems = cronBatchItemRepository.selectFailedByBatchId(batchId);
|
128 |
List<CronBatchItem> failedItems = cronBatchItemRepository.selectFailedByBatchId(batchId);
|
| 91 |
|
129 |
|
| 92 |
int failureCount = failedItems.size();
|
130 |
int failureCount = failedItems.size();
|
| - |
|
131 |
int successCount = (int) cronBatchItemRepository
|
| - |
|
132 |
.countByBatchIdAndStatus(batchId, CronBatchItemStatus.SUCCESS);
|
| 93 |
int successCount = batch.getTotalCount() - failureCount;
|
133 |
int pendingCount = batch.getTotalCount() - successCount - failureCount;
|
| 94 |
|
134 |
|
| 95 |
batch.setSuccessCount(successCount);
|
135 |
batch.setSuccessCount(successCount);
|
| 96 |
batch.setFailureCount(failureCount);
|
136 |
batch.setFailureCount(failureCount);
|
| 97 |
batch.setCompletedAt(LocalDateTime.now());
|
137 |
batch.setCompletedAt(LocalDateTime.now());
|
| - |
|
138 |
if (pendingCount > 0) {
|
| - |
|
139 |
batch.setStatus(CronBatchStatus.INCOMPLETE);
|
| - |
|
140 |
} else if (failureCount > 0) {
|
| 98 |
batch.setStatus(failureCount == 0 ? CronBatchStatus.COMPLETED : CronBatchStatus.PARTIAL_FAILURE);
|
141 |
batch.setStatus(CronBatchStatus.PARTIAL_FAILURE);
|
| - |
|
142 |
} else {
|
| - |
|
143 |
batch.setStatus(CronBatchStatus.COMPLETED);
|
| - |
|
144 |
}
|
| 99 |
|
145 |
|
| 100 |
LOGGER.info("Batch {} finalized: {} success, {} failed", batchId, successCount, failureCount);
|
146 |
LOGGER.info("Batch {} finalized as {}: {} success, {} failed, {} never attempted",
|
| - |
|
147 |
batchId, batch.getStatus(), successCount, failureCount, pendingCount);
|
| 101 |
|
148 |
|
| 102 |
if (failureCount > 0) {
|
149 |
if (failureCount > 0 || pendingCount > 0) {
|
| 103 |
sendFailureEmail(batch, failedItems);
|
150 |
sendFailureEmail(batch, failedItems, pendingCount);
|
| 104 |
}
|
151 |
}
|
| - |
|
152 |
return new BatchOutcome(batch.getStatus(), batch.getTotalCount(),
|
| - |
|
153 |
successCount, failureCount, pendingCount);
|
| 105 |
}
|
154 |
}
|
| 106 |
|
155 |
|
| - |
|
156 |
/**
|
| - |
|
157 |
* Result of a finalize, returned rather than re-read by the caller: finalizeBatch
|
| - |
|
158 |
* commits in its own session, so a caller that re-reads inside its own transaction
|
| - |
|
159 |
* gets the stale pre-finalize entity back out of Hibernate's first-level cache.
|
| - |
|
160 |
*/
|
| - |
|
161 |
public static class BatchOutcome {
|
| - |
|
162 |
private final CronBatchStatus status;
|
| - |
|
163 |
private final int totalCount;
|
| - |
|
164 |
private final int successCount;
|
| - |
|
165 |
private final int failureCount;
|
| - |
|
166 |
private final int neverAttemptedCount;
|
| - |
|
167 |
|
| - |
|
168 |
BatchOutcome(CronBatchStatus status, int totalCount, int successCount,
|
| - |
|
169 |
int failureCount, int neverAttemptedCount) {
|
| - |
|
170 |
this.status = status;
|
| - |
|
171 |
this.totalCount = totalCount;
|
| - |
|
172 |
this.successCount = successCount;
|
| - |
|
173 |
this.failureCount = failureCount;
|
| - |
|
174 |
this.neverAttemptedCount = neverAttemptedCount;
|
| - |
|
175 |
}
|
| - |
|
176 |
|
| - |
|
177 |
public CronBatchStatus getStatus() { return status; }
|
| - |
|
178 |
public int getTotalCount() { return totalCount; }
|
| - |
|
179 |
public int getSuccessCount() { return successCount; }
|
| - |
|
180 |
public int getFailureCount() { return failureCount; }
|
| - |
|
181 |
public int getNeverAttemptedCount() { return neverAttemptedCount; }
|
| - |
|
182 |
}
|
| - |
|
183 |
|
| 107 |
private void sendFailureEmail(CronBatch batch, List<CronBatchItem> failedItems) {
|
184 |
private void sendFailureEmail(CronBatch batch, List<CronBatchItem> failedItems, int pendingCount) {
|
| 108 |
StringBuilder body = new StringBuilder();
|
185 |
StringBuilder body = new StringBuilder();
|
| 109 |
body.append("Cron job: ").append(batch.getJobName()).append("\n");
|
186 |
body.append("Cron job: ").append(batch.getJobName()).append("\n");
|
| - |
|
187 |
body.append("Batch id: ").append(batch.getId()).append("\n");
|
| 110 |
body.append("Run time: ").append(batch.getStartedAt().format(DTF)).append("\n");
|
188 |
body.append("Run time: ").append(batch.getStartedAt().format(DTF)).append("\n");
|
| - |
|
189 |
body.append("Outcome: ").append(batch.getStatus()).append("\n");
|
| 111 |
body.append("Total processed: ").append(batch.getTotalCount()).append("\n");
|
190 |
body.append("Total: ").append(batch.getTotalCount()).append("\n");
|
| 112 |
body.append("Success: ").append(batch.getSuccessCount()).append("\n");
|
191 |
body.append("Success: ").append(batch.getSuccessCount()).append("\n");
|
| 113 |
body.append("Failures: ").append(batch.getFailureCount()).append("\n\n");
|
192 |
body.append("Failed: ").append(batch.getFailureCount()).append("\n");
|
| - |
|
193 |
body.append("Never attempted: ").append(pendingCount).append("\n\n");
|
| - |
|
194 |
|
| - |
|
195 |
if (!failedItems.isEmpty()) {
|
| 114 |
body.append("Partner failures:\n");
|
196 |
body.append("Partner failures:\n");
|
| - |
|
197 |
for (CronBatchItem item : failedItems) {
|
| - |
|
198 |
body.append("- ").append(item.getPartnerName())
|
| - |
|
199 |
.append(" (fofoId: ").append(item.getFofoId()).append(")")
|
| - |
|
200 |
.append(" — ").append(item.getErrorMessage())
|
| - |
|
201 |
.append("\n");
|
| - |
|
202 |
}
|
| - |
|
203 |
body.append("\n");
|
| - |
|
204 |
}
|
| 115 |
|
205 |
|
| - |
|
206 |
if (pendingCount > 0) {
|
| - |
|
207 |
body.append("The run did not reach these ").append(pendingCount)
|
| - |
|
208 |
.append(" item(s) — they were never attempted and no work was done for them:\n");
|
| 116 |
for (CronBatchItem item : failedItems) {
|
209 |
List<CronBatchItem> pendingItems = cronBatchItemRepository
|
| - |
|
210 |
.selectByBatchIdAndStatus(batch.getId(), CronBatchItemStatus.PENDING);
|
| - |
|
211 |
for (CronBatchItem item : pendingItems.subList(0, Math.min(pendingItems.size(), PENDING_LIST_CAP))) {
|
| 117 |
body.append("- ").append(item.getPartnerName())
|
212 |
body.append("- ").append(item.getPartnerName())
|
| 118 |
.append(" (fofoId: ").append(item.getFofoId()).append(")")
|
213 |
.append(" (fofoId: ").append(item.getFofoId()).append(")\n");
|
| - |
|
214 |
}
|
| - |
|
215 |
if (pendingItems.size() > PENDING_LIST_CAP) {
|
| 119 |
.append(" — ").append(item.getErrorMessage())
|
216 |
body.append("... and ").append(pendingItems.size() - PENDING_LIST_CAP).append(" more\n");
|
| - |
|
217 |
}
|
| 120 |
.append("\n");
|
218 |
body.append("\nRe-run the job to process them.\n");
|
| 121 |
}
|
219 |
}
|
| 122 |
|
220 |
|
| - |
|
221 |
String subject = pendingCount > 0
|
| - |
|
222 |
? String.format("[CRON ALERT] %s — INCOMPLETE, %d of %d never attempted",
|
| - |
|
223 |
batch.getJobName(), pendingCount, batch.getTotalCount())
|
| 123 |
String subject = String.format("[CRON ALERT] %s — %d of %d partners failed",
|
224 |
: String.format("[CRON ALERT] %s — %d of %d partners failed",
|
| 124 |
batch.getJobName(), batch.getFailureCount(), batch.getTotalCount());
|
225 |
batch.getJobName(), batch.getFailureCount(), batch.getTotalCount());
|
| 125 |
|
226 |
|
| 126 |
try {
|
227 |
try {
|
| 127 |
mailOutboxService.queueMailViaGoogle(TECHNOLOGY_EMAIL, null, subject, body.toString(),
|
228 |
mailOutboxService.queueMailViaGoogle(TECHNOLOGY_EMAIL, null, subject, body.toString(),
|
| 128 |
"CronBatchService." + batch.getJobName());
|
229 |
"CronBatchService." + batch.getJobName());
|
| 129 |
} catch (Exception e) {
|
230 |
} catch (Exception e) {
|