Implementing Batch Processing
1. Understanding Batch vs Stream Processing
2. Implementing Job Scheduling
| Scheduler | Detail |
|---|---|
| cron / systemd timers | Single host |
| K8s CronJob | Cluster-scheduled |
| Airflow / Dagster / Prefect | DAGs + dependencies |
| Argo Workflows | K8s-native DAGs |
| Quartz | JVM cluster scheduler |
| Temporal | Durable workflow scheduling |
3. Implementing Job Parallelization
| Technique | Detail |
|---|---|
| Data parallelism | Partition by key/range, workers per shard |
| Task parallelism | Independent tasks in DAG |
| MapReduce | map → shuffle → reduce |
| Spark RDD/DataFrame | Stage-based parallel exec |
| K8s parallel Jobs | parallelism + completions |
4. Implementing Checkpointing and Recovery
| Property | Detail |
|---|---|
| Checkpoint | Persist progress (offset, last id) |
| Atomic | Write checkpoint + side effects together (txn) |
| Resume | Restart from last checkpoint |
| Flink savepoint | External, versioned snapshot |
| Spark checkpoint | Truncates lineage; HDFS/S3 |
5. Implementing Error Handling in Batch Jobs
| Pattern | Detail |
|---|---|
| Retry transient | Exponential backoff per task |
| Skip + DLQ | Bad records to side output |
| Fail job | For systemic errors (cred, schema) |
| Idempotent steps | Safe re-run after partial fail |
6. Implementing Batch Partitioning
| Strategy | Detail |
|---|---|
| Range | By id ranges or date partitions |
| Hash | Even distribution |
| Time | Daily / hourly partitions |
| Skew handling | Salting hot keys; Spark AQE skew join |
7. Understanding Batch Size Optimization
| Trade-off | Detail |
|---|---|
| Small | Lower latency, more overhead per record |
| Large | Better throughput, more memory, longer to retry |
| Tune by | Benchmarks, target throughput, memory limits |
| Common | JDBC: 500-5000; Kafka linger 5-50ms |
8. Implementing Job Monitoring
| Metric | Detail |
|---|---|
| Records processed/sec | Throughput |
| Lag | Source - consumed offset |
| Failed records | Count + DLQ depth |
| Job duration | Distribution; alert on outliers |
| Resource use | CPU, mem, executor heap |
9. Implementing Job Retry Strategies
| Level | Detail |
|---|---|
| Task-level | Spark retries failed task N times |
| Job-level | K8s Job backoffLimit |
| Workflow-level | Airflow retries per task |
| Manual rerun | From checkpoint, after fix |
10. Implementing Batch Result Aggregation
| Pattern | Detail |
|---|---|
| Reduce / fold | Sum, count, top-K per partition |
| Tree aggregation | treeReduce for skew tolerance |
| Sketches | HyperLogLog, t-digest for approx |
| Output | Partitioned files (Parquet) or DB upsert |