Skip to content

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

MetricDefinitionTarget
ThroughputEvents processed per second>100K events/sec
Latency (P50)Median end-to-end latency<50ms
Latency (P99)99th percentile latency<200ms
Checkpoint DurationTime to complete checkpoint<5 minutes
State SizeTotal operator state memory< Available RAM
BackpressureIs 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:

FormatSer SpeedDeser SpeedSizeUse Case
JSON50 MB/s40 MB/s100%Human-readable
Bincode500 MB/s800 MB/s30%Internal
Protobuf200 MB/s300 MB/s40%Cross-lang
MessagePack150 MB/s180 MB/s50%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 × Parallelism

Example:

10,000 events/sec × 600 sec window × 1 KB/event × 16 parallelism
= 10,000 × 600 × 1024 × 16 bytes
= 96 GB memory required

3. 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):

Terminal window
# G1GC for large heaps
JAVA_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:

  1. Check backpressure: Is pipeline saturated?
  2. Measure per-operator latency
  3. Review checkpoint duration

Solutions:

  1. Increase parallelism
  2. Optimize expensive operators (profiling)
  3. Reduce checkpoint frequency
  4. Add more resources

Low Throughput

Symptom: Processing < 10K events/sec on powerful hardware

Diagnosis:

  1. Check CPU utilization (should be 70-80%)
  2. Review parallelism settings
  3. Measure serialization overhead
  4. Check for stragglers (slow tasks)

Solutions:

  1. Increase parallelism
  2. Use faster serialization (bincode)
  3. Enable operator fusion
  4. Increase batch size

High Memory Usage

Symptom: OOM errors or excessive GC

Diagnosis:

  1. Measure state size per operator
  2. Check window sizes
  3. Review state TTL configuration

Solutions:

  1. Enable state TTL
  2. Reduce window size
  3. Increase parallelism (distribute state)
  4. Use RocksDB state backend (off-heap)

Backpressure

Symptom: Pipeline slowing down, buffers full

Diagnosis:

  1. Identify bottleneck operator (slowest throughput)
  2. Check resource utilization (CPU, memory, I/O)

Solutions:

  1. Increase parallelism of bottleneck operator
  2. Optimize bottleneck logic
  3. Scale out cluster
  4. Reduce source ingestion rate

Best Practices Summary

  1. Start with defaults, measure, then optimize
  2. Monitor everything: throughput, latency, state size
  3. Test at production scale before deploying
  4. Use batching for throughput-sensitive workloads
  5. Reduce state for latency-sensitive workloads
  6. Set state TTL to prevent unbounded growth
  7. Enable compression for checkpoint storage
  8. Use efficient serialization (bincode, protobuf)
  9. Profile before optimizing (don’t guess!)
  10. Scale horizontally before vertical optimization

References


Version: 1.0