Install
$ agentstack add skill-mickeyyaya-refactoring-skills-batch-job-patterns Open-source listing, not yet scanned by AgentStack. Follow the source repository for install instructions.
Security review
⚠ Flagged1 finding(s); flagged for manual review. · v0.1.0 How review works →
- • Prompt-injection patterns
- • Secret / credential exfiltration
- • Dangerous shell & filesystem operations
- • Untrusted network calls
- • Known-malicious package signatures
- high Dangerous shell/eval execution.
What it can access
- ✓ Network access No
- ● Filesystem access Used
- ✓ Shell / process execution No
- ✓ Environment & secrets No
- ● Dynamic code execution Used
From automated source analysis of v0.1.0. “Used” means the capability is present in the source — more access means more to trust, not that it’s unsafe.
Reliability & compatibility
Declared compatibility
Compatibility is declared by the source manifest. End-to-end runtime verification is coming, see below.
We're building live execution health for every listing: tool-call success rate, median latency, uptime, and last-checked timestamps, measured, not self-reported. It isn't live yet, so we don't show numbers we can't stand behind.
How agent discovery & health will work →About
Batch Job Patterns
Overview
Batch jobs fail in subtle ways: two instances start simultaneously and corrupt shared state, a crash at item 50,000 restarts the job from item 1, a dead worker holds a lock forever, or an unbounded batch runs out of memory. Use this guide when designing, implementing, or reviewing batch processing systems.
When to use: Designing scheduled jobs, ETL pipelines, bulk data migrations, report generation, queue-draining workers, or any process that operates on a bounded or streaming set of records.
Quick Reference
| Pattern | Core Idea | Primary Red Flag | |---------|-----------|-----------------| | Distributed Locking | Only one instance runs at a time via SETNX / advisory lock | Multiple instances starting the same job simultaneously | | Idempotent Checkpoint/Resume | Track cursor position so a crash restarts mid-batch, not from scratch | Job restarts from item 1 on every failure | | Heartbeat / Dead Job Detection | Worker renews a lease; expired lease means worker is dead | Lock held forever by a crashed worker | | Job Scheduling | Cron, interval, event-triggered, or priority queue dispatch | Drift, missed runs, or runaway overlapping executions | | Graceful Shutdown | SIGTERM drains in-flight items before exit | Partial item writes or corrupted state on deploy/restart | | Retry and DLQ | Per-item retry with skip-or-fail policy; unprocessable items route to DLQ | Silent discard of failed items, or one bad item halts the entire batch | | Batch Size Optimization | Memory-bounded chunks with throughput tuning | OOM crashes or 1-item-at-a-time throughput bottlenecks |
Patterns in Detail
1. Distributed Locking for Exclusive Job Execution
Without a distributed lock, two cron nodes or scaled replicas can start the same job simultaneously — leading to duplicate processing, race conditions, and data corruption.
Red Flags:
- No lock before starting a batch job in a multi-instance deployment
- Lock acquired with no expiry — dead worker holds it forever
- Lock released in a non-atomic way (check-then-delete can delete another owner's lock)
- Advisory lock never checked for success before proceeding
Redis SETNX (TypeScript):
// SETNX = SET if Not eXists — atomic exclusive lock
async function acquireLock(
redis: Redis,
lockKey: string,
ttlMs: number,
ownerId: string
): Promise {
// NX: only set if not exists; PX: expire in milliseconds
const result = await redis.set(lockKey, ownerId, 'NX', 'PX', ttlMs);
return result === 'OK';
}
async function releaseLock(redis: Redis, lockKey: string, ownerId: string): Promise {
// Lua script ensures atomic compare-and-delete — only release our own lock
const script = `
if redis.call("get", KEYS[1]) == ARGV[1] then
return redis.call("del", KEYS[1])
else
return 0
end
`;
await redis.eval(script, 1, lockKey, ownerId);
}
async function runExclusiveJob(redis: Redis, jobFn: () => Promise): Promise {
const ownerId = crypto.randomUUID();
const lockKey = 'job:nightly-report:lock';
const acquired = await acquireLock(redis, lockKey, 60_000, ownerId);
if (!acquired) {
logger.info('Job already running on another instance — skipping');
return;
}
try {
await jobFn();
} finally {
await releaseLock(redis, lockKey, ownerId);
}
}
PostgreSQL Advisory Lock (Go):
// Advisory locks are session-scoped and auto-released on connection close
func acquireAdvisoryLock(db *sql.DB, lockID int64) (bool, error) {
var acquired bool
err := db.QueryRow("SELECT pg_try_advisory_lock($1)", lockID).Scan(&acquired)
return acquired, err
}
func runExclusiveJob(db *sql.DB, jobFn func() error) error {
const lockID = 123456789 // unique per job type — use a constant registry
acquired, err := acquireAdvisoryLock(db, lockID)
if err != nil {
return fmt.Errorf("acquireAdvisoryLock: %w", err)
}
if !acquired {
log.Println("job already running — skipping")
return nil
}
defer db.Exec("SELECT pg_advisory_unlock($1)", lockID)
return jobFn()
}
2. Idempotent Checkpoint / Resume Pattern
A job that processes 1M records and crashes at record 750,000 should resume at 750,000 — not restart from 1. Checkpoints also enable idempotency: replaying the job produces the same result.
Red Flags:
- No cursor or offset stored between batches — job always starts from the beginning
- Checkpoint written inside a transaction that also modifies business data, but not atomically
- Resuming without verifying checkpoint integrity (stale or corrupt checkpoint)
- Processing without idempotency keys — duplicate processing on retry
Python (cursor-based pagination with checkpoint):
import json, time
from pathlib import Path
CHECKPOINT_FILE = Path("/var/run/jobs/user-sync.checkpoint.json")
def load_checkpoint() -> dict:
if CHECKPOINT_FILE.exists():
return json.loads(CHECKPOINT_FILE.read_text())
return {"last_id": 0, "processed": 0}
def save_checkpoint(state: dict) -> None:
# Atomic write: write to temp, then rename — prevents partial writes
tmp = CHECKPOINT_FILE.with_suffix(".tmp")
tmp.write_text(json.dumps(state))
tmp.rename(CHECKPOINT_FILE)
def run_user_sync(db, batch_size: int = 500) -> None:
state = load_checkpoint()
last_id = state["last_id"]
total = state["processed"]
while True:
rows = db.query(
"SELECT id, email FROM users WHERE id > %s ORDER BY id LIMIT %s",
(last_id, batch_size)
)
if not rows:
break
for row in rows:
sync_user_to_crm(row) # idempotent — CRM upserts by email
last_id = rows[-1]["id"]
total += len(rows)
save_checkpoint({"last_id": last_id, "processed": total})
CHECKPOINT_FILE.unlink(missing_ok=True) # clean up on success
TypeScript (database-persisted checkpoint):
interface JobCheckpoint {
jobId: string;
cursor: string;
processedCount: number;
updatedAt: Date;
}
async function upsertCheckpoint(db: Db, checkpoint: JobCheckpoint): Promise {
await db.query(
`INSERT INTO job_checkpoints (job_id, cursor, processed_count, updated_at)
VALUES ($1, $2, $3, NOW())
ON CONFLICT (job_id) DO UPDATE
SET cursor = $2, processed_count = $3, updated_at = NOW()`,
[checkpoint.jobId, checkpoint.cursor, checkpoint.processedCount]
);
}
async function resumableBatch(db: Db, jobId: string): Promise {
const saved = await db.query(
'SELECT * FROM job_checkpoints WHERE job_id = $1', [jobId]
);
let cursor = saved.rows[0]?.cursor ?? '';
let processed = saved.rows[0]?.processedCount ?? 0;
while (true) {
const items = await fetchPage(db, cursor, 500);
if (items.length === 0) break;
await processBatch(items);
cursor = items[items.length - 1].id;
processed += items.length;
await upsertCheckpoint(db, { jobId, cursor, processedCount: processed, updatedAt: new Date() });
}
}
3. Heartbeat-Based Dead Job Detection
A worker acquires a lock, then crashes before releasing it. Without a heartbeat mechanism, the lock is held until TTL expires — which may be hours. Heartbeats allow detecting dead workers in seconds.
Red Flags:
- Lock TTL set to job duration — no renewal means lock expires mid-run
- No heartbeat thread — crashed worker's lock persists until TTL
- Heartbeat failure not treated as a fatal error — job continues without a valid lease
- Lease renewal racing with job completion
Go (heartbeat goroutine with lease renewal):
type Lease struct {
redis *redis.Client
key string
ownerID string
ttl time.Duration
stopCh chan struct{}
doneCh chan struct{}
}
func NewLease(redis *redis.Client, key, ownerID string, ttl time.Duration) *Lease {
return &Lease{redis: redis, key: key, ownerID: ownerID, ttl: ttl,
stopCh: make(chan struct{}), doneCh: make(chan struct{})}
}
// StartHeartbeat renews the lease at ttl/2 interval
func (l *Lease) StartHeartbeat(ctx context.Context) {
go func() {
defer close(l.doneCh)
ticker := time.NewTicker(l.ttl / 2)
defer ticker.Stop()
for {
select {
case bool:
last_beat = redis_client.get(HEARTBEAT_KEY.format(job_id=job_id))
if last_beat is None:
return True
elapsed = time.time() - float(last_beat)
return elapsed > DEAD_THRESHOLD
def claim_dead_job(redis_client, job_id: str) -> bool:
"""Claim a stalled job for reprocessing."""
if not is_job_dead(redis_client, job_id):
return False
# Atomic: only one claimer wins the SETNX race
return redis_client.set(
f"job:{job_id}:lock", "new-owner", nx=True, ex=60
)
4. Job Scheduling Strategies
Batch jobs run on a schedule (cron), at fixed intervals, in response to events, or via a priority queue. Choosing the wrong strategy causes drift, missed runs, or starvation.
Red Flags:
- Cron with no overlap protection — previous run still active when next fires
- Interval timer restarts immediately after failure — no backoff
- Event-triggered jobs with no deduplication — fan-out storms on burst events
- Priority queues where low-priority items starve indefinitely
Scheduling strategy comparison:
| Strategy | When to Use | Key Risk | Mitigation | |----------|------------|----------|-----------| | Cron | Fixed wall-clock schedule | Overlap if job > interval | Distributed lock + skip-if-running | | Interval | Run N seconds after last completion | Drift under load | Track completion time, not start time | | Event-triggered | React to upstream data arrival | Duplicate events, fan-out storms | Idempotency key + debounce window | | Priority queue | Mixed urgency workloads | Low-priority starvation | Aging: boost priority after wait threshold |
TypeScript (cron with overlap guard):
import { CronJob } from 'cron';
let running = false;
const job = new CronJob('0 2 * * *', async () => { // 02:00 daily
if (running) {
logger.warn('Previous run still active — skipping this tick');
return;
}
running = true;
try {
await runNightlyReport();
} catch (err) {
logger.error('Nightly report failed', { err });
metrics.increment('batch.nightly_report.failure');
} finally {
running = false;
}
});
Go (priority queue with aging):
type JobItem struct {
ID string
Priority int
EnqueuedAt time.Time
}
// Age items: boost priority by 1 for each 5 minutes of wait
func effectivePriority(item JobItem) int {
ageMinutes := int(time.Since(item.EnqueuedAt).Minutes())
return item.Priority + ageMinutes/5
}
5. Graceful Shutdown During Batch Runs
Deploying or scaling down while a batch is running can leave items in a half-processed state. SIGTERM handling with a drain period ensures the current item finishes before the process exits.
Red Flags:
process.exit()called immediately on SIGTERM — in-flight items are abandoned- No drain timeout — a misbehaving item can prevent shutdown indefinitely
- Database transaction open at shutdown — connection closed mid-transaction causes partial writes
- Checkpoint not saved before exit — resume restarts from last checkpoint, not current cursor
TypeScript (SIGTERM with drain):
let shuttingDown = false;
process.on('SIGTERM', () => {
logger.info('SIGTERM received — draining current batch item');
shuttingDown = true;
// Force exit after drain period regardless
setTimeout(() => {
logger.error('Drain timeout exceeded — forcing exit');
process.exit(1);
}, 30_000).unref();
});
async function processBatchLoop(items: AsyncIterable): Promise {
for await (const item of items) {
if (shuttingDown) {
logger.info('Shutdown flag set — stopping batch loop cleanly');
break;
}
await processItem(item);
await saveCheckpoint(item.id);
}
}
Go (context cancellation on SIGTERM):
func main() {
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer stop()
if err := runBatch(ctx); err != nil && !errors.Is(err, context.Canceled) {
log.Fatalf("batch failed: %v", err)
}
log.Println("batch completed or drained gracefully")
}
func runBatch(ctx context.Context) error {
for {
select {
case bool:
return isinstance(err, (ConnectionError, TimeoutError))
def process_with_retry(
item: dict,
handler: Callable,
dlq_writer,
max_attempts: int = 3,
base_delay: float = 0.5
) -> str:
"""Returns 'ok', 'skipped', or 'dlq'."""
last_err = None
for attempt in range(1, max_attempts + 1):
try:
handler(item)
return 'ok'
except Exception as err:
last_err = err
if not is_retryable(err):
break # permanent failure — skip to DLQ immediately
if attempt BatchResult:
result = BatchResult()
for item in items:
outcome = process_with_retry(item, handler, dlq_writer)
if outcome == 'ok':
result.processed += 1
else:
result.failed += 1
return result
TypeScript (fail-all vs skip policy):
type ItemPolicy = 'skip-on-error' | 'fail-all';
async function processBatch(
items: Item[],
handler: (item: Item) => Promise,
policy: ItemPolicy = 'skip-on-error',
dlq: DLQClient
): Promise {
const errors: Array = [];
for (const item of items) {
try {
await handler(item);
} catch (err) {
if (policy === 'fail-all') throw err; // abort entire batch
logger.warn('Item failed — routing to DLQ', { itemId: item.id, err });
await dlq.send({ item, error: String(err), ts: new Date().toISOString() });
errors.push({ item, err });
}
}
if (errors.length > 0) {
logger.error(`Batch completed with ${errors.length} item failures routed to DLQ`);
metrics.gauge('batch.dlq_routed', errors.length);
}
}
Cross-reference: error-handling-patterns — Dead Letter Queue for DLQ monitoring and replay patterns.
7. Batch Size Optimization
Too small: context-switch overhead dominates throughput. Too large: OOM crash or lock contention. Optimal batch size is bounded by memory and tuned for throughput.
Red Flags:
- Batch size hardcoded to 1 — N round-trips instead of 1 bulk insert
- Batch size unbounded —
SELECT * FROM tableloads all rows into memory - No memory accounting — batch size set by row count, not by actual byte footprint
- No throughput metrics — batch size never tuned based on observed performance
Memory-bounded batching (Go):
const (
MaxBatchBytes = 64 * 1024 * 1024 // 64 MB cap per batch
MinBatchSize = 10
MaxBatchSize = 5_000
)
func memoryBoundedBatch(rows []Row, estimateSize func(Row) int) [][]Row {
var batches [][]Row
var current []Row
var currentBytes int
for _, row := range rows {
rowBytes := estimateSize(row)
if len(current) >= MinBatchSize &&
(currentBytes+rowBytes > MaxBatchBytes || len(current) >= MaxBatchSize) {
batches = append(batches, current)
current = nil
currentBytes = 0
}
current = append(current, row)
currentBytes += rowBytes
}
if len(current) > 0 {
batches = append(batches, current)
}
return batches
}
Throughput tuning with adaptive batch size (Python):
import time
class AdaptiveBatcher:
"""Adjusts batch size based on observed throughput."""
def __init__(self, initial_size: int = 100, min_size: int = 10, max_size: int = 2000):
self.size = initial_size
self.min_size = min_size
self.max_size = max_size
self._last_throughput: float | None = None
…
## Source & license
This open-source skill is cataloged on AgentStack and links to its original source — we do not rehost the code.
- **Author:** [mickeyyaya](https://github.com/mickeyyaya)
- **Source:** [mickeyyaya/refactoring-skills](https://github.com/mickeyyaya/refactoring-skills)
- **License:** MIT
Install and usage instructions live in the source repository linked above.
Reviews
No reviews yet, be the first.
Write a review
Versions
- v0.1.0 Imported from the upstream source.