Exactly-Once Semantics Guide
Exactly-Once Semantics Guide
Overview
HeliosDB Streaming guarantees exactly-once processing semantics using a combination of:
- Distributed Checkpointing with AES-256-GCM encryption
- Two-Phase Commit (2PC) for sinks
- Idempotent Source reads with offset tracking
- 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:
- Detect Failure: Operator fails or becomes unresponsive
- Stop Processing: All operators stop accepting new events
- Load Checkpoint: Retrieve last successful checkpoint from storage
- Decrypt State: Decrypt operator state using key management
- Restore State: Restore all operator state to checkpoint
- Resume Sources: Sources resume from checkpointed offsets
- Abort Uncommitted: Sinks abort any pre-committed but uncommitted transactions
- 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:
| Configuration | Checkpoint Duration | Throughput Impact | Latency Impact |
|---|---|---|---|
| No encryption | 2.3s | <1% | +5ms P99 |
| AES-256-GCM | 2.8s | <2% | +8ms P99 |
| With compression | 2.1s | <1% | +6ms P99 |
Recommendation: Use encryption + compression for best security/performance balance.
2PC Latency Impact
Per-batch latency overhead:
| Batch Size | 2PC Overhead | Total Latency |
|---|---|---|
| 100 events | +12ms | 45ms |
| 1000 events | +18ms | 62ms |
| 10000 events | +35ms | 110ms |
Recommendation: Batch 1000-5000 events for optimal throughput/latency.
State Size Impact
Encrypted checkpoint size vs raw state:
| Raw State Size | Encrypted Size | Compression Ratio |
|---|---|---|
| 100 MB | 105 MB | 1.05x |
| 1 GB | 1.08 GB | 1.08x |
| 10 GB | 10.2 GB | 1.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:
- Increase
timeoutin checkpoint config - Reduce state size by enabling TTL
- Use faster checkpoint storage (SSD vs HDD)
High Checkpoint Duration
Symptom: Checkpoints taking >5 minutes
Solutions:
- Enable compression
- Increase checkpoint interval to reduce frequency
- Scale out parallelism to reduce per-task state
Recovery Failures
Symptom: Job fails to recover from checkpoint
Solutions:
- Verify checkpoint storage is accessible
- Check encryption keys are available
- Validate checkpoint format compatibility
- Review logs for specific error messages
References
- Flink Checkpointing Documentation
- HeliosDB Encryption: CHECKPOINT_ENCRYPTION.md
Version: 1.0