返回博客

历史数据回放:时间旅行查询与状态重建

coomia-dip 的历史回放系统基于 Iceberg 的时间旅行和 Nessie 的分支版本管理,实现了任意时间点的数据状态查询和完整状态重建。系统支持快照查询("某一时刻数据什么样")、增量回放("两个时间点之间发生了什么变化")和假设分析("如果那个变更没发生,现在数据会是什么样")。本文从时间旅行架构、版本管理机制、回放引擎实现到生产最佳实践,完整解析这一企业级历史数据能力。

Coomia发布于 2025年9月24日10 分钟阅读
分享本文Twitter / X

系列:S6 平台工程 · 第 10 篇 | 难度:高级 | 阅读时间:18 分钟

历史数据回放:时间旅行查询与状态重建

#TL;DR

coomia-dip 的历史回放系统基于 Iceberg 的时间旅行和 Nessie 的分支版本管理,实现了任意时间点的数据状态查询和完整状态重建。系统支持快照查询("某一时刻数据什么样")、增量回放("两个时间点之间发生了什么变化")和假设分析("如果那个变更没发生,现在数据会是什么样")。本文从时间旅行架构、版本管理机制、回放引擎实现到生产最佳实践,完整解析这一企业级历史数据能力。

#1. 历史回放的核心价值

#1.1 业务场景

  • 合规审计:监管要求查看某一时点的数据状态
  • 故障恢复:数据被误修改后回滚到正确状态
  • 根因分析:追溯数据异常的起始时间和变更链
  • 假设分析:模拟"如果当时没有执行某个操作"的场景

#1.2 技术基础

coomia-dip 的历史回放建立在两个核心技术之上:

Code
┌─────────────────────────────────────────────┐
│              History Replay Engine            │
│  ┌─────────────────┐  ┌──────────────────┐  │
│  │ Apache Iceberg   │  │ Project Nessie   │  │
│  │ (时间旅行快照)    │  │ (分支版本管理)    │  │
│  │                  │  │                  │  │
│  │ • 快照隔离       │  │ • Git-like 分支  │  │
│  │ • 时间旅行查询   │  │ • 原子提交       │  │
│  │ • 模式演进       │  │ • 合并/冲突解决   │  │
│  └─────────────────┘  └──────────────────┘  │
└─────────────────────────────────────────────┘

#1.3 对标 Palantir Foundry

能力Palantir Foundrycoomia-dip
时间旅行Transaction 日志Iceberg 快照
分支管理有限Nessie Git-like 分支
增量查询支持Iceberg 增量读
假设分析部分支持Nessie 分支 + 回放
保留策略不透明可配置快照保留

#2. 时间旅行架构

#2.1 快照管理

Python
class SnapshotManager:
    """Iceberg 快照管理器"""

    async def query_at_timestamp(
        self,
        table_name: str,
        timestamp: datetime,
        filter_expr: str | None = None,
    ) -> ResultSet:
        """查询指定时间点的数据状态"""
        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:
        """查询指定快照的数据状态"""
        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]:
        """列出时间范围内的所有快照"""
        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 增量变更查询

Python
class IncrementalChangeQuery:
    """增量变更查询 - 获取两个时间点之间的变更"""

    async def get_changes_between(
        self,
        table_name: str,
        start_time: datetime,
        end_time: datetime,
    ) -> ChangeSet:
        """获取两个时间点之间的所有变更"""
        table = await self._catalog.load_table(table_name)

        start_snapshot = self._find_snapshot(table, start_time)
        end_snapshot = self._find_snapshot(table, end_time)

        # 使用 Iceberg 增量读取
        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]:
        """获取单个字段的变更历史"""
        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 分支管理

#3.1 分支操作

Python
class NessieBranchManager:
    """Nessie 分支管理器 - Git-like 数据版本控制"""

    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 假设分析

Python
class WhatIfAnalyzer:
    """假设分析引擎 - 基于 Nessie 分支"""

    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. 回放引擎

#4.1 状态重建

Python
class StateRebuilder:
    """状态重建引擎 - 从历史数据重建完整状态"""

    async def rebuild_object_state(
        self,
        object_type: str,
        object_id: str,
        target_time: datetime,
    ) -> ObjectState:
        """重建指定时间点的对象状态"""
        # 1. 从 Iceberg 获取时间旅行数据
        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,
            )

        # 2. 重建对象(包括关联关系)
        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 变更回放

Python
class ChangeReplayer:
    """变更回放引擎"""

    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 服务接口

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. 测试策略

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. 生产最佳实践

#7.1 快照保留策略

数据重要性快照频率保留时间归档策略
核心业务每次写入365 天归档到对象存储
一般业务每小时90 天压缩后归档
临时数据每天30 天不归档

#7.2 性能考虑

  • 时间旅行查询性能与快照数量相关,建议定期合并小快照
  • 增量查询远快于全量比较,优先使用
  • 假设分析分支应设置自动过期时间
  • 大规模回滚前必须执行 dry_run 评估影响

#7.3 安全考虑

  • 回滚操作需要特殊权限(ADMIN 或 DATA_STEWARD)
  • 假设分析分支的数据不应暴露给普通用户
  • 历史查询同样受 ABAC 和脱敏策略控制

#8. 总结

coomia-dip 的历史回放系统通过 Iceberg 时间旅行和 Nessie 分支管理的组合,提供了企业级的历史数据能力。关键设计亮点:

  1. 时间旅行:任意时间点的数据快照查询
  2. 增量变更:高效获取两个时间点之间的数据差异
  3. 假设分析:基于分支的 what-if 场景模拟
  4. 安全回滚:dry-run 预评估 + 受控回滚执行
  5. 版本控制:Git-like 数据分支管理

下一篇将深入探讨 coomia-dip 的元数据目录系统。