Implementing Batch Processing

1. Understanding Batch vs Stream Processing

Batch

  • Bounded data; periodic
  • High throughput, high latency
  • Spark, MapReduce, Flink batch

Stream

  • Unbounded; continuous
  • Low latency, per-event
  • Flink, Kafka Streams, Spark Structured Streaming

2. Implementing Job Scheduling

SchedulerDetail
cron / systemd timersSingle host
K8s CronJobCluster-scheduled
Airflow / Dagster / PrefectDAGs + dependencies
Argo WorkflowsK8s-native DAGs
QuartzJVM cluster scheduler
TemporalDurable workflow scheduling

3. Implementing Job Parallelization

TechniqueDetail
Data parallelismPartition by key/range, workers per shard
Task parallelismIndependent tasks in DAG
MapReducemap → shuffle → reduce
Spark RDD/DataFrameStage-based parallel exec
K8s parallel Jobsparallelism + completions

4. Implementing Checkpointing and Recovery

PropertyDetail
CheckpointPersist progress (offset, last id)
AtomicWrite checkpoint + side effects together (txn)
ResumeRestart from last checkpoint
Flink savepointExternal, versioned snapshot
Spark checkpointTruncates lineage; HDFS/S3

5. Implementing Error Handling in Batch Jobs

PatternDetail
Retry transientExponential backoff per task
Skip + DLQBad records to side output
Fail jobFor systemic errors (cred, schema)
Idempotent stepsSafe re-run after partial fail

6. Implementing Batch Partitioning

StrategyDetail
RangeBy id ranges or date partitions
HashEven distribution
TimeDaily / hourly partitions
Skew handlingSalting hot keys; Spark AQE skew join

7. Understanding Batch Size Optimization

Trade-offDetail
SmallLower latency, more overhead per record
LargeBetter throughput, more memory, longer to retry
Tune byBenchmarks, target throughput, memory limits
CommonJDBC: 500-5000; Kafka linger 5-50ms

8. Implementing Job Monitoring

MetricDetail
Records processed/secThroughput
LagSource - consumed offset
Failed recordsCount + DLQ depth
Job durationDistribution; alert on outliers
Resource useCPU, mem, executor heap

9. Implementing Job Retry Strategies

LevelDetail
Task-levelSpark retries failed task N times
Job-levelK8s Job backoffLimit
Workflow-levelAirflow retries per task
Manual rerunFrom checkpoint, after fix

10. Implementing Batch Result Aggregation

PatternDetail
Reduce / foldSum, count, top-K per partition
Tree aggregationtreeReduce for skew tolerance
SketchesHyperLogLog, t-digest for approx
OutputPartitioned files (Parquet) or DB upsert