feat: retry-механизм (max 3 попытки)
- V2__add_retry_count.sql: retry_count int NOT NULL DEFAULT 0 - takeNextJob: new + failed (retry_count < 3) + stale processing (30с) - markFailed: retry_count = retry_count + 1 - JobProcessor: логирование retry (X/3) - Heartbeat: не трогать (уже работает, 10с)
This commit is contained in:
@@ -26,6 +26,7 @@ data class Job(
|
|||||||
val createdAt: Instant,
|
val createdAt: Instant,
|
||||||
val updatedAt: Instant,
|
val updatedAt: Instant,
|
||||||
val takenAt: Instant? = null,
|
val takenAt: Instant? = null,
|
||||||
|
val retryCount: Int = 0,
|
||||||
) {
|
) {
|
||||||
companion object {
|
companion object {
|
||||||
const val STATUS_NEW = "new"
|
const val STATUS_NEW = "new"
|
||||||
@@ -33,5 +34,6 @@ data class Job(
|
|||||||
const val STATUS_DONE = "done"
|
const val STATUS_DONE = "done"
|
||||||
const val STATUS_FAILED = "failed"
|
const val STATUS_FAILED = "failed"
|
||||||
const val STATUS_CANCELLED = "cancelled"
|
const val STATUS_CANCELLED = "cancelled"
|
||||||
|
const val MAX_RETRIES = 3
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -71,7 +71,7 @@ open class JobRepository(
|
|||||||
jdbcTemplate.update(
|
jdbcTemplate.update(
|
||||||
"""
|
"""
|
||||||
UPDATE media_mirror.jobs
|
UPDATE media_mirror.jobs
|
||||||
SET status = 'failed', error = ?, updated_at = now()
|
SET status = 'failed', error = ?, retry_count = retry_count + 1, updated_at = now()
|
||||||
WHERE id = ?
|
WHERE id = ?
|
||||||
""".trimIndent(),
|
""".trimIndent(),
|
||||||
error.take(4000),
|
error.take(4000),
|
||||||
@@ -110,8 +110,16 @@ open class JobRepository(
|
|||||||
val jobs = jdbcTemplate.query(
|
val jobs = jdbcTemplate.query(
|
||||||
"""
|
"""
|
||||||
SELECT * FROM media_mirror.jobs
|
SELECT * FROM media_mirror.jobs
|
||||||
WHERE status = 'new'
|
WHERE (
|
||||||
ORDER BY created_at
|
status = 'new'
|
||||||
|
OR (status = 'failed' AND retry_count < 3)
|
||||||
|
OR (status = 'processing' AND taken_at < now() - interval '30 seconds')
|
||||||
|
)
|
||||||
|
ORDER BY
|
||||||
|
CASE WHEN status = 'new' THEN 0
|
||||||
|
WHEN status = 'failed' THEN 1
|
||||||
|
ELSE 2 END,
|
||||||
|
created_at
|
||||||
LIMIT 1
|
LIMIT 1
|
||||||
FOR UPDATE SKIP LOCKED
|
FOR UPDATE SKIP LOCKED
|
||||||
""".trimIndent(),
|
""".trimIndent(),
|
||||||
@@ -152,6 +160,7 @@ open class JobRepository(
|
|||||||
createdAt = rs.getTimestamp("created_at").toInstant(),
|
createdAt = rs.getTimestamp("created_at").toInstant(),
|
||||||
updatedAt = rs.getTimestamp("updated_at").toInstant(),
|
updatedAt = rs.getTimestamp("updated_at").toInstant(),
|
||||||
takenAt = rs.getTimestamp("taken_at")?.toInstant(),
|
takenAt = rs.getTimestamp("taken_at")?.toInstant(),
|
||||||
|
retryCount = rs.getInt("retry_count"),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -58,6 +58,12 @@ class JobProcessor(
|
|||||||
throw e
|
throw e
|
||||||
} catch (e: Throwable) {
|
} catch (e: Throwable) {
|
||||||
logger.warn("Job {} finished with error: {}", job.id, e.message)
|
logger.warn("Job {} finished with error: {}", job.id, e.message)
|
||||||
|
val currentRetries = jobRepository.findById(job.id)?.retryCount ?: 0
|
||||||
|
if (currentRetries >= Job.MAX_RETRIES) {
|
||||||
|
logger.warn("Job {} exceeded max retries ({}), permanently failed", job.id, currentRetries)
|
||||||
|
} else {
|
||||||
|
logger.info("Job {} failed (retry {}/{}), will be retried", job.id, currentRetries, Job.MAX_RETRIES)
|
||||||
|
}
|
||||||
jobRepository.markFailed(job.id, e.message ?: e.toString())
|
jobRepository.markFailed(job.id, e.message ?: e.toString())
|
||||||
} finally {
|
} finally {
|
||||||
heartbeatJob.cancel()
|
heartbeatJob.cancel()
|
||||||
|
|||||||
@@ -0,0 +1 @@
|
|||||||
|
ALTER TABLE media_mirror.jobs ADD COLUMN IF NOT EXISTS retry_count int NOT NULL DEFAULT 0;
|
||||||
Reference in New Issue
Block a user