Architecture Overview
When building real-time image processing systems, traditional request-response patterns often can't meet the dual demands of high throughput and low latency. We need an event-driven pipeline architecture that splits image processing into multiple independent stages, each independently scalable and optimizable.
Core design principles: decouple, scale, fault-tolerant. Each processing stage is an independent service communicating through message queues, supporting horizontal scaling and fault isolation.
A pipeline's value lies not in the speed of any single stage, but in the coordination and balance of the whole. One bottleneck can slow down the entire chain.
Pipeline Stage Design
We split image processing into 5 core stages, each corresponding to an independent microservice:
- Ingest — Receives raw images, performs format validation and metadata extraction, supporting HTTP / gRPC / WebSocket access methods
- Preprocess — Image scaling, color space conversion, normalization, data augmentation, outputting uniformly formatted tensors
- Inference — Model inference with batch processing and dynamic batching support, loading multiple models (classification, detection, segmentation)
- Postprocess — Result parsing, NMS, confidence filtering, coordinate mapping back to original image
- Export — Result persistence to databases, callback notifications, Webhook pushes, message queue broadcasts
Each stage connects through Apache Kafka, supporting backpressure mechanisms to prevent downstream overload. Kafka's partitioning allows us to partition by image ID, ensuring processing order for the same image.
Kafka Streaming Architecture
Kafka serves as the pipeline's core message middleware, handling data flow, backpressure control, and fault recovery. We designed the following Topic structure:
- image.ingest — Raw image metadata, 3 partitions, 7-day retention
- image.preprocessed — Preprocessed tensor data, 6 partitions
- image.inference.result — Inference results, 12 partitions
- image.export — Final processing results, 6 partitions
Key design decisions:
- Message Size — Image data isn't placed directly in Kafka; it's stored in S3/MinIO, with Kafka messages containing only reference pointers
- Consumer Groups — Each stage uses an independent consumer group, supporting independent scaling
- Dead Letter Queue — Failed messages automatically route to DLQ, preventing pipeline blockage
Lesson learned: smaller Kafka messages enable higher throughput. A proven approach is to keep large payloads, such as images, in object storage and pass only metadata and references through Kafka.
Worker Concurrency Model
Each pipeline stage's Worker is implemented in Go, leveraging goroutines for high-concurrency processing. The core uses the Worker Pool pattern with channels for task distribution and result collection.
Worker Pool design essentials:
- Fixed-Size Pool — Prevents unbounded goroutine growth causing OOM; pool size dynamically adjusts based on CPU cores and GPU memory
- Graceful Shutdown — Propagates cancellation signals via context, ensuring no tasks are lost during rolling updates
- Health Checks — Each Worker exposes a /health endpoint; Kubernetes probes detect liveness status
- Metrics Collection — Uses Prometheus client to record processing latency, throughput, and error rates
Kubernetes Deployment
Each pipeline stage runs as an independent Kubernetes Deployment. An HPA (Horizontal Pod Autoscaler) scales it according to queue depth. Here is the deployment configuration for the Inference Worker:
Key configuration: minReplicas: 2, maxReplicas: 20, targetQueueLength: 100. When queue depth exceeds the threshold, HPA completes scaling within 30 seconds. GPU nodes use NVIDIA Device Plugin and node affinity for scheduling to dedicated GPU node pools.
Monitoring & Observability
We adopted the three pillars of observability, ensuring every pipeline stage is transparent and controllable:
- Metrics — Prometheus collects throughput, P99 latency, error rates per stage; Grafana Dashboard displays in real-time
- Logging — Structured logs collected via Fluentd to Elasticsearch, supporting real-time queries and alerting rules
- Tracing — OpenTelemetry implements end-to-end distributed tracing, visualizing the full chain from image ingestion to result export
Key alerting rules:
- Trigger P1 alert when Kafka Consumer Lag exceeds 1000
- Trigger auto-scaling when inference latency P99 exceeds 500ms
- Notify development team when Dead Letter Queue message count grows continuously
- Trigger scale-down when GPU utilization drops below 30%, saving resource costs
Scaling Strategies
In production, we adopted a layered scaling strategy, ensuring efficient operation under different loads:
- Horizontal Scaling — Each stage has independent HPA, scaling based on its own load without mutual impact
- Vertical Scaling — GPU-intensive stages use node affinity to schedule on dedicated nodes, reserving GPU resources
- Geographic Scaling — Multi-region deployment for proximity access, reducing network latency; Kafka MirrorMaker enables cross-region data sync
- Elastic Scaling — Warm-up + rapid scaling strategy for traffic spikes; KEDA enables precise scaling based on Kafka lag
The final system runs stably in production, processing 10M+ images daily, with P99 latency under 200ms and 99.99% availability. During peak events, it successfully handled 5x traffic spikes with auto-scaling completing within 2 minutes.
Summary & Lessons
Building a real-time image processing pipeline is a systems engineering effort spanning architecture design, message middleware selection, container orchestration, and monitoring. Our core lessons:
- Decouple First, Optimize Later — Prioritize architectural scalability early; optimize bottleneck stages later
- Separate Data from Control — Large data goes to object storage, control flow to message queues, preventing Kafka from becoming a bottleneck
- Observability First — Deploy complete monitoring before going live; otherwise you can't diagnose issues
- Fault-Tolerant Design — Every stage should assume downstream failure, using DLQ, retries, and circuit breakers for system resilience
This architecture has been running stably for over a year, withstanding multiple business peaks and infrastructure changes. If you're building a similar system, I hope these experiences provide useful reference.