Back to Blog

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.

CoomiaPublished on September 24, 20258 min read
Share this articleTwitter / X

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:

Code
┌─────────────────────────────────────────────┐
│              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

CapabilityPalantir Foundrycoomia-dip
Time travelTransaction logIceberg snapshots
Branch managementLimitedNessie Git-like branching
Incremental queriesSupportedIceberg incremental reads
What-if analysisPartial supportNessie branches + replay
Retention policiesOpaqueConfigurable snapshot retention

#2. Time-Travel Architecture

#2.1 Snapshot Management

Python
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

Python
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

Python
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

Python
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

Python
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

Python
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

PROTOBUF
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

Python
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 ImportanceSnapshot FrequencyRetentionArchive Strategy
Core businessEvery write365 daysArchive to object storage
General businessHourly90 daysCompressed archive
Temporary dataDaily30 daysNo 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:

  1. Time travel: Snapshot queries at any point in time
  2. Incremental changes: Efficiently retrieve data differences between two timestamps
  3. What-if analysis: Branch-based scenario simulation
  4. Safe rollback: Dry-run pre-assessment + controlled rollback execution
  5. Version control: Git-like data branch management

The next article will explore the coomia-dip metadata catalog system.