Skip to content

Exactly-Once Semantics Guide

Exactly-Once Semantics Guide

Overview

HeliosDB Streaming guarantees exactly-once processing semantics using a combination of:

  1. Distributed Checkpointing with AES-256-GCM encryption
  2. Two-Phase Commit (2PC) for sinks
  3. Idempotent Source reads with offset tracking
  4. State Recovery from encrypted checkpoints

This ensures that each event is processed exactly once, even in the presence of failures.

How It Works

Architecture

Source → Operators → Checkpoint Barrier → 2PC Sink → External System
↓ ↓ ↓
Offset Track State Snapshot Pre-commit → Commit
↓
Encrypted Storage
(S3/Distributed FS)

1. Checkpointing Mechanism

HeliosDB periodically creates consistent snapshots of all operator state.

2. Checkpoint Barriers

Checkpoint barriers flow through the dataflow graph ensuring all operators checkpoint at the same logical time.

Barrier Alignment:

Stream A: [e1, e2, e3, |barrier-42|, e4, e5, ...]
Stream B: [e6, e7, |barrier-42|, e8, e9, e10, ...]
↓
Aligned Checkpoint
(both streams at barrier-42)

3. Two-Phase Commit Sinks

Sinks implement 2PC to ensure atomic writes to external systems.

4. Recovery Semantics

On failure, the system recovers from the last successful checkpoint.

Recovery Process:

  1. Detect Failure: Operator fails or becomes unresponsive
  2. Stop Processing: All operators stop accepting new events
  3. Load Checkpoint: Retrieve last successful checkpoint from storage
  4. Decrypt State: Decrypt operator state using key management
  5. Restore State: Restore all operator state to checkpoint
  6. Resume Sources: Sources resume from checkpointed offsets
  7. Abort Uncommitted: Sinks abort any pre-committed but uncommitted transactions
  8. Resume Processing: Processing continues from checkpoint

Exactly-Once Guarantees

What Is Guaranteed

Each event processed exactly once within the streaming application State updates are consistent across all operators Sinks write each result exactly once to external systems No data loss on failures No duplicate writes to external systems

What Is NOT Guaranteed

❌ External system reads may see partial results during recovery ❌ Side effects in user functions (e.g., logging, metrics) may occur multiple times ❌ Non-transactional sinks cannot provide exactly-once (only at-least-once)

Performance Impact

Checkpoint Overhead

Measured on a representative workload (10K events/sec); figures are indicative, reproduce on your own hardware:

ConfigurationCheckpoint DurationThroughput ImpactLatency Impact
No encryption2.3s<1%+5ms P99
AES-256-GCM2.8s<2%+8ms P99
With compression2.1s<1%+6ms P99

Recommendation: Use encryption + compression for best security/performance balance.

2PC Latency Impact

Per-batch latency overhead:

Batch Size2PC OverheadTotal Latency
100 events+12ms45ms
1000 events+18ms62ms
10000 events+35ms110ms

Recommendation: Batch 1000-5000 events for optimal throughput/latency.

State Size Impact

Encrypted checkpoint size vs raw state:

Raw State SizeEncrypted SizeCompression Ratio
100 MB105 MB1.05x
1 GB1.08 GB1.08x
10 GB10.2 GB1.02x

Conclusion: Encryption overhead is minimal (< 10%).

Best Practices

1. Checkpoint Interval Tuning

Very frequent checkpoints (for example every 10 seconds) add overhead and reduce throughput; very infrequent checkpoints (for example every 10 minutes) lengthen recovery after a failure. An interval of around 60 seconds, with a minimum pause between checkpoints to prevent thrashing, balances recovery time and overhead.

2. Monitor Checkpoint Health

Track checkpoint duration, state size and failures, and alert when a checkpoint takes unusually long (for example more than 5 minutes).

3. Test Recovery Scenarios Regularly

Inject failures in test environments and verify that results remain exactly-once after recovery.

4. Use Idempotent Sinks When Possible

Even if 2PC fails, idempotent sinks prevent duplicates, for example by writing with an upsert (INSERT ... ON CONFLICT (id) DO UPDATE).

Troubleshooting

Checkpoint Timeouts

Symptom: Checkpoints failing with timeout errors

Solutions:

  1. Increase timeout in checkpoint config
  2. Reduce state size by enabling TTL
  3. Use faster checkpoint storage (SSD vs HDD)

High Checkpoint Duration

Symptom: Checkpoints taking >5 minutes

Solutions:

  1. Enable compression
  2. Increase checkpoint interval to reduce frequency
  3. Scale out parallelism to reduce per-task state

Recovery Failures

Symptom: Job fails to recover from checkpoint

Solutions:

  1. Verify checkpoint storage is accessible
  2. Check encryption keys are available
  3. Validate checkpoint format compatibility
  4. Review logs for specific error messages

References


Version: 1.0