Performance Tuning Guide
Performance Tuning Guide
Overview
This guide covers performance optimization techniques for HeliosDB Streaming applications. Topics include throughput optimization, latency reduction, memory management, and monitoring setup.
Performance Dimensions
Key Metrics
| Metric | Definition | Target |
|---|---|---|
| Throughput | Events processed per second | >100K events/sec |
| Latency (P50) | Median end-to-end latency | <50ms |
| Latency (P99) | 99th percentile latency | <200ms |
| Checkpoint Duration | Time to complete checkpoint | <5 minutes |
| State Size | Total operator state memory | < Available RAM |
| Backpressure | Is pipeline saturated? | No backpressure |
Throughput Optimization
1. Parallelism Configuration
Set appropriate parallelism based on CPU cores and data volume.
Guidelines:
- Start with parallelism = CPU cores
- Increase for I/O-bound operations (2-3x cores)
- Decrease for memory-intensive operations
- Monitor CPU utilization (target 70-80%)
2. Batching
Batch events to amortize per-event overhead.
Trade-offs:
- Higher throughput (10-50x improvement)
- Lower CPU overhead
- ❌ Higher latency (wait for full batch)
- ❌ More memory usage
3. Buffer Sizing
Configure internal buffers to prevent bottlenecks.
Guidelines:
- Small buffers (1K): Low latency, risk of backpressure
- Medium buffers (10K): Balanced
- Large buffers (100K): High throughput, more memory
4. Serialization Optimization
Use efficient serialization formats.
Format Comparison:
| Format | Ser Speed | Deser Speed | Size | Use Case |
|---|---|---|---|---|
| JSON | 50 MB/s | 40 MB/s | 100% | Human-readable |
| Bincode | 500 MB/s | 800 MB/s | 30% | Internal |
| Protobuf | 200 MB/s | 300 MB/s | 40% | Cross-lang |
| MessagePack | 150 MB/s | 180 MB/s | 50% | Compact |
5. Operator Fusion
Fuse operators to reduce overhead.
Improvement: 50-70% throughput increase for simple pipelines.
Latency Optimization
1. Reduce Checkpoint Frequency
Trade recovery time for lower latency.
Frequent checkpoints (for example every 10 seconds) add roughly 5-10ms to P99 latency; infrequent checkpoints (for example every 5 minutes) add roughly 1-2ms.
2. Minimize State
Stateless operators have lower latency.
Stateful operators (for example keyed aggregations) require state access, which adds latency; stateless operators such as pure map/filter functions avoid it.
3. Optimize Serialization
Use zero-copy serialization where possible.
4. Window Size Reduction
Smaller windows trigger more frequently (lower latency).
Memory Management
1. State TTL
Set time-to-live for state to prevent unbounded growth.
2. Window Sizing
Balance window size with memory usage.
Memory Calculation:
Memory = Events_per_sec × Window_duration_sec × Event_size × ParallelismExample:
10,000 events/sec × 600 sec window × 1 KB/event × 16 parallelism= 10,000 × 600 × 1024 × 16 bytes= 96 GB memory required3. State Backend Configuration
Choose appropriate state backend: in-memory (fastest, limited size) or RocksDB (larger state, slower).
4. Garbage Collection Tuning
For JVM-based operators (rare in Rust, but if using JNI):
# G1GC for large heapsJAVA_OPTS="-XX:+UseG1GC -XX:MaxGCPauseMillis=50 -Xmx32g"Monitoring & Observability
1. Prometheus Metrics
Export key metrics for monitoring: streaming_events_processed_total (throughput), streaming_latency_ms (end-to-end latency), streaming_state_size_bytes (state size) and streaming_backpressure (0 or 1).
2. Grafana Dashboard
Dashboard JSON Template:
{ "dashboard": { "title": "HeliosDB Streaming Performance", "panels": [ { "title": "Throughput (events/sec)", "targets": [ { "expr": "rate(streaming_events_processed_total[1m])" } ], "type": "graph" }, { "title": "Latency (P50, P95, P99)", "targets": [ { "expr": "histogram_quantile(0.50, streaming_latency_ms)", "legendFormat": "P50" }, { "expr": "histogram_quantile(0.95, streaming_latency_ms)", "legendFormat": "P95" }, { "expr": "histogram_quantile(0.99, streaming_latency_ms)", "legendFormat": "P99" } ], "type": "graph" }, { "title": "State Size (MB)", "targets": [ { "expr": "streaming_state_size_bytes / 1024 / 1024" } ], "type": "graph" }, { "title": "Checkpoint Duration (seconds)", "targets": [ { "expr": "streaming_checkpoint_duration_seconds" } ], "type": "graph" }, { "title": "Backpressure Status", "targets": [ { "expr": "streaming_backpressure" } ], "type": "stat" } ] }}3. Logging Best Practices
use tracing::{info, warn, error, debug, instrument};
#[instrument(skip(event), fields(event_id = %event.id))]async fn process_event(event: Event) -> Result<ProcessedEvent> { debug!("Starting event processing");
let start = Instant::now(); let result = expensive_operation(&event).await?; let duration = start.elapsed();
if duration > Duration::from_millis(100) { warn!("Slow processing: {:?}", duration); }
info!("Event processed successfully"); Ok(result)}4. Alert Rules
Prometheus AlertManager Rules:
groups: - name: streaming_alerts rules: - alert: HighLatency expr: histogram_quantile(0.99, streaming_latency_ms) > 500 for: 5m labels: severity: warning annotations: summary: "High P99 latency detected" description: "P99 latency is {{ $value }}ms (threshold: 500ms)"
- alert: Backpressure expr: streaming_backpressure > 0 for: 2m labels: severity: critical annotations: summary: "Pipeline experiencing backpressure"
- alert: CheckpointFailure expr: increase(streaming_checkpoint_failures_total[5m]) > 2 labels: severity: critical annotations: summary: "Multiple checkpoint failures"Troubleshooting
High Latency
Symptom: P99 latency > 500ms
Diagnosis:
- Check backpressure: Is pipeline saturated?
- Measure per-operator latency
- Review checkpoint duration
Solutions:
- Increase parallelism
- Optimize expensive operators (profiling)
- Reduce checkpoint frequency
- Add more resources
Low Throughput
Symptom: Processing < 10K events/sec on powerful hardware
Diagnosis:
- Check CPU utilization (should be 70-80%)
- Review parallelism settings
- Measure serialization overhead
- Check for stragglers (slow tasks)
Solutions:
- Increase parallelism
- Use faster serialization (bincode)
- Enable operator fusion
- Increase batch size
High Memory Usage
Symptom: OOM errors or excessive GC
Diagnosis:
- Measure state size per operator
- Check window sizes
- Review state TTL configuration
Solutions:
- Enable state TTL
- Reduce window size
- Increase parallelism (distribute state)
- Use RocksDB state backend (off-heap)
Backpressure
Symptom: Pipeline slowing down, buffers full
Diagnosis:
- Identify bottleneck operator (slowest throughput)
- Check resource utilization (CPU, memory, I/O)
Solutions:
- Increase parallelism of bottleneck operator
- Optimize bottleneck logic
- Scale out cluster
- Reduce source ingestion rate
Best Practices Summary
- Start with defaults, measure, then optimize
- Monitor everything: throughput, latency, state size
- Test at production scale before deploying
- Use batching for throughput-sensitive workloads
- Reduce state for latency-sensitive workloads
- Set state TTL to prevent unbounded growth
- Enable compression for checkpoint storage
- Use efficient serialization (bincode, protobuf)
- Profile before optimizing (don’t guess!)
- Scale horizontally before vertical optimization
References
Version: 1.0