Historical Replay: Time-Travel Queries and State Reconstruction
The coomia-dip historical replay system leverages Iceberg's time-travel capabilities and Nessie's branch-based version management to enable data state queries at any point in time and complete state reconstruction. The system supports snapshot queries ("what did the data look like at a given moment"), incremental replay ("what changes occurred between two points in time"), and what-if analysis ("what would the data look like if that change hadn't happened"). This article covers the complete architecture from time-travel design, version management, replay engine implementation, to production best practices.
“Series: S6 Platform Engineering · Article 10 | Level: Advanced | Reading Time: 18 min
Historical Replay: Time-Travel Queries and State Reconstruction
#TL;DR
The coomia-dip historical replay system leverages Iceberg's time-travel capabilities and Nessie's branch-based version management to enable data state queries at any point in time and complete state reconstruction. The system supports snapshot queries ("what did the data look like at a given moment"), incremental replay ("what changes occurred between two points in time"), and what-if analysis ("what would the data look like if that change hadn't happened"). This article covers the complete architecture from time-travel design, version management, replay engine implementation, to production best practices.
#1. Core Value of Historical Replay
#1.1 Business Scenarios
- Compliance auditing: Regulators require viewing data state at specific points in time
- Fault recovery: Roll back to correct state after accidental data modifications
- Root cause analysis: Trace the start time and change chain of data anomalies
- What-if analysis: Simulate "what if that operation hadn't been executed" scenarios
#1.2 Technical Foundation
coomia-dip historical replay is built on two core technologies:
┌─────────────────────────────────────────────┐
│ History Replay Engine │
│ ┌─────────────────┐ ┌──────────────────┐ │
│ │ Apache Iceberg │ │ Project Nessie │ │
│ │ (Time-Travel │ │ (Branch Version │ │
│ │ Snapshots) │ │ Management) │ │
│ │ │ │ │ │
│ │ • Snapshot │ │ • Git-like │ │
│ │ isolation │ │ branching │ │
│ │ • Time-travel │ │ • Atomic commits │ │
│ │ queries │ │ • Merge/conflict │ │
│ │ • Schema │ │ resolution │ │
│ │ evolution │ │ │ │
│ └─────────────────┘ └──────────────────┘ │
└─────────────────────────────────────────────┘
#1.3 Comparison with Palantir Foundry
| Capability | Palantir Foundry | coomia-dip |
|---|---|---|
| Time travel | Transaction log | Iceberg snapshots |
| Branch management | Limited | Nessie Git-like branching |
| Incremental queries | Supported | Iceberg incremental reads |
| What-if analysis | Partial support | Nessie branches + replay |
| Retention policies | Opaque | Configurable snapshot retention |
#2. Time-Travel Architecture
#2.1 Snapshot Management
class SnapshotManager:
"""Iceberg Snapshot Manager"""
async def query_at_timestamp(
self,
table_name: str,
timestamp: datetime,
filter_expr: str | None = None,
) -> ResultSet:
"""Query data state at a specific point in time"""
table = await self._catalog.load_table(table_name)
target_snapshot = None
for snapshot in table.snapshots():
snap_time = datetime.fromtimestamp(snapshot.timestamp_ms / 1000)
if snap_time <= timestamp:
target_snapshot = snapshot
else:
break
if target_snapshot is None:
raise HistoryNotAvailable(f"No snapshot found before {timestamp}")
scan = table.scan(snapshot_id=target_snapshot.snapshot_id)
if filter_expr:
scan = scan.filter(filter_expr)
return await self._execute_scan(scan)
async def query_at_snapshot(
self, table_name: str, snapshot_id: int,
) -> ResultSet:
"""Query data state at a specific snapshot"""
table = await self._catalog.load_table(table_name)
scan = table.scan(snapshot_id=snapshot_id)
return await self._execute_scan(scan)
async def list_snapshots(
self,
table_name: str,
start_time: datetime | None = None,
end_time: datetime | None = None,
) -> list[SnapshotInfo]:
"""List all snapshots within a time range"""
table = await self._catalog.load_table(table_name)
snapshots = []
for snapshot in table.snapshots():
snap_time = datetime.fromtimestamp(snapshot.timestamp_ms / 1000)
if start_time and snap_time < start_time:
continue
if end_time and snap_time > end_time:
continue
snapshots.append(SnapshotInfo(
snapshot_id=snapshot.snapshot_id,
timestamp=snap_time,
operation=snapshot.operation,
summary=snapshot.summary,
))
return snapshots
#2.2 Incremental Change Queries
class IncrementalChangeQuery:
"""Incremental change query - get changes between two points in time"""
async def get_changes_between(
self, table_name: str, start_time: datetime, end_time: datetime,
) -> ChangeSet:
"""Get all changes between two timestamps"""
table = await self._catalog.load_table(table_name)
start_snapshot = self._find_snapshot(table, start_time)
end_snapshot = self._find_snapshot(table, end_time)
added_records = []
deleted_records = []
for data_file in end_snapshot.added_files_since(start_snapshot):
records = await self._read_data_file(data_file)
added_records.extend(records)
for data_file in end_snapshot.deleted_files_since(start_snapshot):
records = await self._read_data_file(data_file)
deleted_records.extend(records)
return ChangeSet(
table_name=table_name,
start_time=start_time,
end_time=end_time,
added=added_records,
deleted=deleted_records,
modified=self._detect_modifications(added_records, deleted_records),
)
async def get_field_history(
self,
table_name: str,
record_id: str,
field_name: str,
start_time: datetime | None = None,
end_time: datetime | None = None,
) -> list[FieldChange]:
"""Get the change history of a single field"""
changes = []
snapshots = await self._snapshot_manager.list_snapshots(
table_name, start_time, end_time,
)
prev_value = None
for snapshot in snapshots:
current = await self._get_field_at_snapshot(
table_name, record_id, field_name, snapshot.snapshot_id,
)
if current != prev_value:
changes.append(FieldChange(
timestamp=snapshot.timestamp,
old_value=prev_value,
new_value=current,
snapshot_id=snapshot.snapshot_id,
))
prev_value = current
return changes
#3. Nessie Branch Management
#3.1 Branch Operations
class NessieBranchManager:
"""Nessie branch manager - Git-like data version control"""
async def create_branch(
self,
branch_name: str,
from_ref: str = "main",
from_timestamp: datetime | None = None,
) -> BranchInfo:
if from_timestamp:
ref = await self._nessie_client.get_reference_at(from_ref, from_timestamp)
else:
ref = await self._nessie_client.get_reference(from_ref)
branch = await self._nessie_client.create_reference(
branch_name=branch_name,
source_hash=ref.hash,
)
return BranchInfo(
name=branch.name,
hash=branch.hash,
created_at=datetime.utcnow(),
source_ref=from_ref,
)
async def merge_branch(
self,
source_branch: str,
target_branch: str = "main",
conflict_resolution: ConflictResolution = ConflictResolution.FAIL,
) -> MergeResult:
source = await self._nessie_client.get_reference(source_branch)
target = await self._nessie_client.get_reference(target_branch)
try:
result = await self._nessie_client.merge(
from_ref=source.name,
to_ref=target.name,
from_hash=source.hash,
)
return MergeResult(success=True, merged_hash=result.hash, conflicts=[])
except NessieConflictError as e:
if conflict_resolution == ConflictResolution.FAIL:
raise
return await self._resolve_conflicts(e.conflicts, conflict_resolution)
async def diff_branches(self, branch_a: str, branch_b: str) -> BranchDiff:
diff = await self._nessie_client.get_diff(branch_a, branch_b)
return BranchDiff(
branch_a=branch_a,
branch_b=branch_b,
added_tables=diff.added,
modified_tables=diff.modified,
deleted_tables=diff.deleted,
)
#3.2 What-If Analysis
class WhatIfAnalyzer:
"""What-If analysis engine based on Nessie branches"""
async def create_what_if_scenario(
self, name: str, base_timestamp: datetime, description: str = "",
) -> WhatIfScenario:
branch = await self._branch_manager.create_branch(
branch_name=f"what-if/{name}",
from_ref="main",
from_timestamp=base_timestamp,
)
return WhatIfScenario(
name=name,
branch_name=branch.name,
base_timestamp=base_timestamp,
description=description,
created_at=datetime.utcnow(),
)
async def apply_hypothetical_changes(
self, scenario: WhatIfScenario, changes: list[HypotheticalChange],
) -> None:
for change in changes:
if change.type == "skip_action":
await self._revert_action_on_branch(
scenario.branch_name, change.action_id,
)
elif change.type == "modify_value":
await self._modify_on_branch(
scenario.branch_name,
change.table_name,
change.record_id,
change.field_name,
change.new_value,
)
async def compare_scenario_with_reality(
self, scenario: WhatIfScenario,
) -> ScenarioComparison:
diff = await self._branch_manager.diff_branches(
scenario.branch_name, "main",
)
return ScenarioComparison(
scenario=scenario,
diff=diff,
divergence_point=scenario.base_timestamp,
)
#4. Replay Engine
#4.1 State Reconstruction
class StateRebuilder:
"""State reconstruction engine"""
async def rebuild_object_state(
self, object_type: str, object_id: str, target_time: datetime,
) -> ObjectState:
table_name = self._get_table_name(object_type)
result = await self._snapshot_manager.query_at_timestamp(
table_name=table_name,
timestamp=target_time,
filter_expr=f"id = '{object_id}'",
)
if not result.rows:
return ObjectState(
object_type=object_type, object_id=object_id,
exists=False, timestamp=target_time,
)
properties = result.rows[0]
links = await self._rebuild_links(object_type, object_id, target_time)
return ObjectState(
object_type=object_type, object_id=object_id,
exists=True, timestamp=target_time,
properties=properties, links=links,
)
async def rebuild_full_state(
self, object_type: str, target_time: datetime,
filter_expr: str | None = None,
) -> list[ObjectState]:
table_name = self._get_table_name(object_type)
result = await self._snapshot_manager.query_at_timestamp(
table_name=table_name,
timestamp=target_time,
filter_expr=filter_expr,
)
states = []
for row in result.rows:
states.append(ObjectState(
object_type=object_type, object_id=row["id"],
exists=True, timestamp=target_time, properties=row,
))
return states
#4.2 Change Replay
class ChangeReplayer:
"""Change replay engine"""
async def replay_changes(
self,
table_name: str,
start_time: datetime,
end_time: datetime,
speed: float = 1.0,
callback: Callable[[ChangeEvent], Awaitable[None]] | None = None,
) -> ReplayResult:
changes = await self._change_query.get_changes_between(
table_name, start_time, end_time,
)
total_changes = len(changes.added) + len(changes.deleted) + len(changes.modified)
processed = 0
for change in changes.in_chronological_order():
if callback:
await callback(change)
processed += 1
if speed < 1.0:
await asyncio.sleep((1.0 / speed) * 0.01)
return ReplayResult(
table_name=table_name,
time_range=(start_time, end_time),
total_changes=total_changes,
processed=processed,
)
async def rollback_to_timestamp(
self, table_name: str, target_time: datetime, dry_run: bool = True,
) -> RollbackPlan:
current_snapshot = await self._snapshot_manager.get_current_snapshot(table_name)
target_snapshot = self._find_snapshot_at(table_name, target_time)
changes_to_revert = await self._change_query.get_changes_between(
table_name, target_time, datetime.utcnow(),
)
plan = RollbackPlan(
table_name=table_name,
current_snapshot=current_snapshot,
target_snapshot=target_snapshot,
records_to_revert=len(changes_to_revert.modified),
records_to_restore=len(changes_to_revert.deleted),
records_to_remove=len(changes_to_revert.added),
)
if not dry_run:
await self._execute_rollback(plan)
return plan
#5. gRPC Service Interface
syntax = "proto3";
package onto.history.v1;
service HistoryService {
rpc QueryAtTimestamp(QueryAtTimestampRequest) returns (QueryResponse);
rpc QueryAtSnapshot(QueryAtSnapshotRequest) returns (QueryResponse);
rpc GetChangesBetween(GetChangesRequest) returns (ChangeSetResponse);
rpc GetFieldHistory(GetFieldHistoryRequest) returns (FieldHistoryResponse);
rpc CreateBranch(CreateBranchRequest) returns (BranchInfo);
rpc MergeBranch(MergeBranchRequest) returns (MergeResult);
rpc DiffBranches(DiffBranchesRequest) returns (BranchDiff);
rpc CreateWhatIfScenario(CreateScenarioRequest) returns (WhatIfScenario);
rpc CompareScenario(CompareScenarioRequest) returns (ScenarioComparison);
rpc RollbackToTimestamp(RollbackRequest) returns (RollbackPlan);
rpc ReplayChanges(ReplayRequest) returns (stream ChangeEvent);
rpc ListSnapshots(ListSnapshotsRequest) returns (SnapshotListResponse);
}
#6. Testing Strategy
class TestHistoricalReplay:
async def test_time_travel_query(self):
await service.create_object("Employee", {"id": "e1", "salary": 50000})
t1 = datetime.utcnow()
await service.update_object("Employee", "e1", {"salary": 60000})
t2 = datetime.utcnow()
state_t1 = await history.query_at_timestamp("employees", t1)
state_t2 = await history.query_at_timestamp("employees", t2)
assert state_t1.rows[0]["salary"] == 50000
assert state_t2.rows[0]["salary"] == 60000
async def test_field_history(self):
changes = await history.get_field_history("employees", "e1", "salary")
assert len(changes) == 2
assert changes[0].new_value == 50000
assert changes[1].new_value == 60000
async def test_what_if_scenario(self):
scenario = await analyzer.create_what_if_scenario(
name="no-raise", base_timestamp=before_raise_time,
)
comparison = await analyzer.compare_scenario_with_reality(scenario)
assert comparison.diff.modified_tables
async def test_rollback_dry_run(self):
plan = await replayer.rollback_to_timestamp(
"employees", one_hour_ago, dry_run=True,
)
assert plan.records_to_revert > 0
#7. Production Best Practices
#7.1 Snapshot Retention Policies
| Data Importance | Snapshot Frequency | Retention | Archive Strategy |
|---|---|---|---|
| Core business | Every write | 365 days | Archive to object storage |
| General business | Hourly | 90 days | Compressed archive |
| Temporary data | Daily | 30 days | No archiving |
#7.2 Performance Considerations
- Time-travel query performance correlates with snapshot count; periodically compact small snapshots
- Incremental queries are much faster than full comparisons; prefer them
- What-if analysis branches should have automatic expiration
- Always run dry_run before large-scale rollbacks to evaluate impact
#7.3 Security Considerations
- Rollback operations require special permissions (ADMIN or DATA_STEWARD)
- What-if analysis branch data should not be exposed to regular users
- Historical queries are subject to the same ABAC and masking policies
#8. Summary
The coomia-dip historical replay system provides enterprise-grade historical data capabilities through the combination of Iceberg time-travel and Nessie branch management. Key design highlights:
- Time travel: Snapshot queries at any point in time
- Incremental changes: Efficiently retrieve data differences between two timestamps
- What-if analysis: Branch-based scenario simulation
- Safe rollback: Dry-run pre-assessment + controlled rollback execution
- Version control: Git-like data branch management
The next article will explore the coomia-dip metadata catalog system.