| Line 24... |
Line 24... |
| 24 |
@Component
|
24 |
@Component
|
| 25 |
public class BatchScheduledTasks {
|
25 |
public class BatchScheduledTasks {
|
| 26 |
|
26 |
|
| 27 |
private static final Logger LOGGER = LogManager.getLogger(BatchScheduledTasks.class);
|
27 |
private static final Logger LOGGER = LogManager.getLogger(BatchScheduledTasks.class);
|
| 28 |
|
28 |
|
| - |
|
29 |
/**
|
| - |
|
30 |
* A RUNNING batch older than this was abandoned by a run that died before finalizing.
|
| - |
|
31 |
* Generous on purpose: the longest real batch (offer processing over ~1,500 partners)
|
| - |
|
32 |
* is minutes, so 6h can only catch a genuinely dead run.
|
| - |
|
33 |
*/
|
| - |
|
34 |
private static final int STALE_BATCH_HOURS = 6;
|
| - |
|
35 |
|
| 29 |
@Autowired
|
36 |
@Autowired
|
| 30 |
private CronBatchService cronBatchService;
|
37 |
private CronBatchService cronBatchService;
|
| 31 |
|
38 |
|
| 32 |
@Autowired
|
39 |
@Autowired
|
| 33 |
private OfferBatchService offerBatchService;
|
40 |
private OfferBatchService offerBatchService;
|
| Line 76... |
Line 83... |
| 76 |
LOGGER.error("Aged stock daily refresh failed: {}", e.getMessage(), e);
|
83 |
LOGGER.error("Aged stock daily refresh failed: {}", e.getMessage(), e);
|
| 77 |
}
|
84 |
}
|
| 78 |
}
|
85 |
}
|
| 79 |
|
86 |
|
| 80 |
/**
|
87 |
/**
|
| - |
|
88 |
* Closes out batches left RUNNING by a run that died mid-loop (the cron JVM is
|
| - |
|
89 |
* kernel-killed on this box often enough that this is routine, not exotic).
|
| - |
|
90 |
*
|
| - |
|
91 |
* Until this existed such a batch stayed RUNNING forever, which mattered twice over:
|
| - |
|
92 |
* it never sent its failure email, and — for offer batches — hasUnfinishedBatch()
|
| - |
|
93 |
* refused every later run, so the job could never be retried. Two were stuck in prod
|
| - |
|
94 |
* for months on that path.
|
| - |
|
95 |
*
|
| - |
|
96 |
* The loop lives HERE and not inside CronBatchService on purpose: a self-invoked
|
| - |
|
97 |
* finalizeBatch() would bypass the Spring proxy, losing REQUIRES_NEW and running with
|
| - |
|
98 |
* no transaction at all, so the writes would be silently discarded.
|
| - |
|
99 |
*/
|
| - |
|
100 |
public void reapStaleBatches() {
|
| - |
|
101 |
List<Integer> staleBatchIds = cronBatchService.selectStaleRunningBatchIds(STALE_BATCH_HOURS);
|
| - |
|
102 |
if (staleBatchIds.isEmpty()) {
|
| - |
|
103 |
return;
|
| - |
|
104 |
}
|
| - |
|
105 |
LOGGER.warn("Reaping {} stale RUNNING batch(es) older than {}h: {}",
|
| - |
|
106 |
staleBatchIds.size(), STALE_BATCH_HOURS, staleBatchIds);
|
| - |
|
107 |
for (Integer batchId : staleBatchIds) {
|
| - |
|
108 |
try {
|
| - |
|
109 |
cronBatchService.finalizeBatch(batchId);
|
| - |
|
110 |
} catch (Exception e) {
|
| - |
|
111 |
LOGGER.error("Failed to reap stale batch {}: {}", batchId, e.getMessage(), e);
|
| - |
|
112 |
}
|
| - |
|
113 |
}
|
| - |
|
114 |
}
|
| - |
|
115 |
|
| - |
|
116 |
/**
|
| 81 |
* CLI entrypoint for cron: delegates to shared OfferBatchService (also used by web/fofo controllers).
|
117 |
* CLI entrypoint for cron: delegates to shared OfferBatchService (also used by web/fofo controllers).
|
| 82 |
*/
|
118 |
*/
|
| 83 |
public void processOfferWithBatch(int offerId) throws Exception {
|
119 |
public void processOfferWithBatch(int offerId) throws Exception {
|
| 84 |
offerBatchService.processOfferWithBatch(offerId);
|
120 |
offerBatchService.processOfferWithBatch(offerId);
|
| 85 |
}
|
121 |
}
|