返回博客

Diff 查询:分支对比与变更追踪

Tags: #DiffQuery #BranchDiff #ChangeTracking #Nessie #Audit #智策平台

Coomia发布于 2025年7月20日14 分钟阅读
分享本文Twitter / X

系列:S3 数据基座 · 第 11 篇 | 难度:高级 | 阅读时间:20 分钟

Diff 查询:分支对比与变更追踪

Tags: #DiffQuery #BranchDiff #ChangeTracking #Nessie #Audit #智策平台

#TL;DR

Diff 查询是 coomia-dip 平台时间旅行能力的进阶应用——不仅能查看历史状态,还能精确定位"什么发生了变化"。本文深入 Nessie 分支模型下的 Diff 查询实现,涵盖多分支对比的语义定义、三种 Diff 算法(全量对比、增量文件对比、Iceberg Changelog 对比)、Schema 演进场景下的兼容 Diff、Diff 结果的结构化输出与可视化、合并冲突检测与解决策略,以及基于 Diff 的变更订阅和审计合规应用。

#1. Diff 查询的场景定义

#1.1 Diff 查询类型

Code
coomia-dip Diff 查询类型矩阵:

┌─────────────────┬───────────────────┬────────────────────┐
│ Diff 类型         │ 对比维度            │ OQL 语法             │
├─────────────────┼───────────────────┼────────────────────┤
│ 时间 Diff        │ 同分支、不同时间点  │ DIFF Entity          │
│                  │                    │ FROM TIME t1         │
│                  │                    │ TO TIME t2           │
├─────────────────┼───────────────────┼────────────────────┤
│ 分支 Diff        │ 不同分支、同一时间  │ DIFF Entity          │
│                  │                    │ FROM BRANCH b1       │
│                  │                    │ TO BRANCH b2         │
├─────────────────┼───────────────────┼────────────────────┤
│ 混合 Diff        │ 不同分支+不同时间  │ DIFF Entity          │
│                  │                    │ FROM BRANCH b1       │
│                  │                    │   AT TIME t1         │
│                  │                    │ TO BRANCH b2         │
│                  │                    │   AT TIME t2         │
├─────────────────┼───────────────────┼────────────────────┤
│ 实体 Diff        │ 单个实体的变更历史  │ DIFF Entity          │
│                  │                    │ WHERE id = 'x'       │
│                  │                    │ FROM TIME t1         │
│                  │                    │ TO TIME t2           │
└─────────────────┴───────────────────┴────────────────────┘

#1.2 Nessie 分支模型回顾

Code
Nessie 分支与 Diff 的关系:

main ──●──●──●──●──●──●──●──●──●── HEAD
        \           \
         \           \──●──●──● feature-b
          \
           \──●──●──●──●── feature-a

Diff 操作:
  1. main vs feature-a → 分支 Diff(自分叉点起的变更)
  2. feature-a vs feature-b → 跨分支 Diff
  3. main@t1 vs main@t2 → 时间 Diff(同分支不同时间)
  4. merge(feature-a, main) → 合并前冲突检测 Diff

#2. Diff 算法实现

#2.1 算法一:全量对比

Python
class FullScanDiffAlgorithm:
    """全量扫描对比算法 — 最简单但最慢"""

    async def diff(
        self,
        entity_type: str,
        from_ref: DataRef,
        to_ref: DataRef,
        columns: list[str] | None = None
    ) -> DiffResult:
        # 读取两个版本的完整数据
        from_data = await self._read_full(entity_type, from_ref)
        to_data = await self._read_full(entity_type, to_ref)

        # 构建哈希索引
        from_map: dict[str, dict] = {}
        for row in from_data:
            from_map[row['entity_id']] = row

        to_map: dict[str, dict] = {}
        for row in to_data:
            to_map[row['entity_id']] = row

        changes = []

        # 检测新增和修改
        for eid, to_row in to_map.items():
            if eid not in from_map:
                changes.append(Change('ADDED', eid, None, to_row))
            else:
                from_row = from_map[eid]
                diff_cols = self._diff_columns(from_row, to_row, columns)
                if diff_cols:
                    changes.append(Change('MODIFIED', eid,
                        {c: from_row[c] for c in diff_cols},
                        {c: to_row[c] for c in diff_cols}
                    ))

        # 检测删除
        for eid in from_map:
            if eid not in to_map:
                changes.append(Change('DELETED', eid, from_map[eid], None))

        return DiffResult(changes=changes)
Code
全量对比算法特性:
  时间复杂度:O(N + M),N 和 M 为两个版本的行数
  空间复杂度:O(N + M)
  优点:简单、正确
  缺点:需要读取全部数据,数据量大时慢

适用场景:
  - 数据量 < 100K 行
  - 首次 Diff(无历史增量信息)

#2.2 算法二:增量文件对比

Python
class IncrementalFileDiffAlgorithm:
    """增量文件对比算法 — 利用 Iceberg 快照差异"""

    async def diff(
        self,
        entity_type: str,
        from_snapshot: int,
        to_snapshot: int,
        columns: list[str] | None = None
    ) -> DiffResult:
        table = self._catalog.load_table(entity_type)

        # 获取两个快照之间的文件差异
        added_files = []
        deleted_files = []

        for entry in table.inspect.entries(
            from_snapshot_id=from_snapshot,
            to_snapshot_id=to_snapshot
        ):
            if entry.status == 'ADDED':
                added_files.append(entry.data_file)
            elif entry.status == 'DELETED':
                deleted_files.append(entry.data_file)

        # 只读取变更的文件(而非全表)
        added_data = self._read_files(added_files, columns)
        deleted_data = self._read_files(deleted_files, columns)

        # 匹配变更
        added_index = {r['entity_id']: r for r in added_data}
        deleted_index = {r['entity_id']: r for r in deleted_data}

        changes = []

        # 在两个集合中都出现 → 修改
        for eid in set(added_index) & set(deleted_index):
            old_row = deleted_index[eid]
            new_row = added_index[eid]
            diff_cols = self._diff_columns(old_row, new_row, columns)
            if diff_cols:
                changes.append(Change('MODIFIED', eid,
                    {c: old_row[c] for c in diff_cols},
                    {c: new_row[c] for c in diff_cols}
                ))

        # 只在 added 中出现 → 新增
        for eid in set(added_index) - set(deleted_index):
            changes.append(Change('ADDED', eid, None, added_index[eid]))

        # 只在 deleted 中出现 → 删除
        for eid in set(deleted_index) - set(added_index):
            changes.append(Change('DELETED', eid, deleted_index[eid], None))

        return DiffResult(changes=changes)
Code
增量文件对比算法特性:
  时间复杂度:O(C),C 为变更的行数
  空间复杂度:O(C)
  优点:只读取变更部分,大数据集高效
  缺点:依赖 Iceberg 文件级追踪

适用场景:
  - 数据量 > 100K 行
  - 变更量 < 总量的 10%
  - 同一个 Iceberg 表的两个快照对比

#2.3 算法三:Changelog 对比

Python
class ChangelogDiffAlgorithm:
    """基于 Iceberg Changelog 的对比算法"""

    async def diff(
        self,
        entity_type: str,
        from_snapshot: int,
        to_snapshot: int
    ) -> DiffResult:
        table = self._catalog.load_table(entity_type)

        # 使用 Iceberg 的原生 changelog API
        changelog = table.inspect.changelog(
            start_snapshot_id=from_snapshot,
            end_snapshot_id=to_snapshot
        )

        changes = []
        for entry in changelog:
            if entry.operation == 'INSERT':
                changes.append(Change('ADDED', entry.record['entity_id'],
                    None, entry.record))
            elif entry.operation == 'DELETE':
                changes.append(Change('DELETED', entry.record['entity_id'],
                    entry.record, None))
            elif entry.operation == 'UPDATE':
                changes.append(Change('MODIFIED', entry.record['entity_id'],
                    entry.before, entry.after))

        return DiffResult(changes=changes)

#3. Schema 演进下的兼容 Diff

#3.1 Schema 变更类型

Code
Schema 演进对 Diff 的影响:

┌─────────────────┬────────────────────────────────────┐
│ Schema 变更       │ Diff 处理策略                        │
├─────────────────┼────────────────────────────────────┤
│ 新增列            │ 新版本有值、旧版本填充 NULL            │
│ 删除列            │ 旧版本有值、新版本标记为 DROPPED       │
│ 列重命名          │ 通过 Iceberg field-id 映射           │
│ 类型变更          │ 自动类型提升(int→long, float→double)│
│ 新增 Entity Type │ 全部标记为 ADDED                     │
│ 删除 Entity Type │ 全部标记为 DELETED                   │
└─────────────────┴────────────────────────────────────┘

#3.2 兼容 Diff 实现

Python
class SchemaAwareDiffEngine:
    """Schema 感知的 Diff 引擎"""

    def diff_with_schema_evolution(
        self,
        from_data: pa.Table,
        to_data: pa.Table,
        from_schema: IcebergSchema,
        to_schema: IcebergSchema
    ) -> DiffResult:
        # 统一 Schema:合并两个版本的所有列
        unified_schema = self._unify_schemas(from_schema, to_schema)

        # 将两个版本的数据投影到统一 Schema
        from_projected = self._project_to_schema(from_data, unified_schema)
        to_projected = self._project_to_schema(to_data, unified_schema)

        # 使用 field-id 而非列名匹配
        column_mapping = self._build_field_id_mapping(
            from_schema, to_schema
        )

        changes = []
        for eid, from_row, to_row in self._aligned_rows(
            from_projected, to_projected
        ):
            schema_changes = []
            value_changes = []

            for col in unified_schema.columns:
                from_val = from_row.get(col.field_id)
                to_val = to_row.get(col.field_id)

                if col.field_id in from_schema and col.field_id not in to_schema:
                    schema_changes.append(SchemaChange(
                        'COLUMN_DROPPED', col.name
                    ))
                elif col.field_id not in from_schema and col.field_id in to_schema:
                    schema_changes.append(SchemaChange(
                        'COLUMN_ADDED', col.name
                    ))
                elif from_val != to_val:
                    value_changes.append(ValueChange(
                        col.name, from_val, to_val
                    ))

            if schema_changes or value_changes:
                changes.append(EnrichedChange(
                    entity_id=eid,
                    value_changes=value_changes,
                    schema_changes=schema_changes
                ))

        return DiffResult(changes=changes)

#4. 合并冲突检测

#4.1 三方 Diff

Code
三方 Diff 用于合并冲突检测:

       base (共同祖先)
      /              \
     /                \
branch-a              branch-b
     \                /
      \              /
       merge (合并点)

三方 Diff 逻辑:
  对于每个 entity_id:
    base_val  = base 分支的值
    a_val     = branch-a 的值
    b_val     = branch-b 的值

    情况 1:base = a = b        → 无变更
    情况 2:base ≠ a, base = b  → a 修改了,取 a
    情况 3:base = a, base ≠ b  → b 修改了,取 b
    情况 4:base ≠ a, base ≠ b, a = b → 双方改成一样,取 a
    情况 5:base ≠ a, base ≠ b, a ≠ b → 冲突!

#4.2 冲突检测实现

Python
class ThreeWayDiff:
    """三方 Diff 冲突检测"""

    def detect_conflicts(
        self,
        base_data: dict[str, dict],
        branch_a_data: dict[str, dict],
        branch_b_data: dict[str, dict]
    ) -> MergeResult:
        all_ids = set(base_data) | set(branch_a_data) | set(branch_b_data)

        auto_resolved = []
        conflicts = []

        for eid in all_ids:
            base = base_data.get(eid)
            a = branch_a_data.get(eid)
            b = branch_b_data.get(eid)

            if base == a == b:
                continue  # 无变更

            if base == a and base != b:
                auto_resolved.append(Resolution(eid, 'TAKE_B', b))
            elif base != a and base == b:
                auto_resolved.append(Resolution(eid, 'TAKE_A', a))
            elif a == b:
                auto_resolved.append(Resolution(eid, 'TAKE_BOTH', a))
            else:
                # 属性级冲突检测
                prop_conflicts = self._detect_property_conflicts(
                    base, a, b
                )
                if prop_conflicts:
                    conflicts.append(Conflict(
                        entity_id=eid,
                        base_value=base,
                        branch_a_value=a,
                        branch_b_value=b,
                        conflicting_properties=prop_conflicts
                    ))
                else:
                    # 不同属性被修改,可自动合并
                    merged = self._auto_merge_properties(base, a, b)
                    auto_resolved.append(Resolution(eid, 'AUTO_MERGE', merged))

        return MergeResult(
            auto_resolved=auto_resolved,
            conflicts=conflicts,
            has_conflicts=len(conflicts) > 0
        )

    def _detect_property_conflicts(
        self, base: dict, a: dict, b: dict
    ) -> list[str]:
        """检测属性级冲突"""
        conflicting = []
        all_keys = set(base or {}) | set(a or {}) | set(b or {})

        for key in all_keys:
            base_val = (base or {}).get(key)
            a_val = (a or {}).get(key)
            b_val = (b or {}).get(key)

            if base_val != a_val and base_val != b_val and a_val != b_val:
                conflicting.append(key)

        return conflicting

#4.3 冲突解决策略

Code
冲突解决策略:

┌──────────────────┬──────────────────────────────────┐
│ 策略               │ 描述                               │
├──────────────────┼──────────────────────────────────┤
│ TAKE_SOURCE       │ 以源分支(被合并方)为准              │
│ TAKE_TARGET       │ 以目标分支(合并到)为准              │
│ TAKE_LATEST       │ 取最后修改时间较新的版本              │
│ MANUAL            │ 人工选择                            │
│ PROPERTY_LEVEL    │ 属性级别自动合并                     │
│                   │ (不同属性分别取各自修改方)            │
│ CUSTOM_FUNCTION   │ 自定义合并函数                       │
└──────────────────┴──────────────────────────────────┘

#5. Diff 结果结构化输出

#5.1 输出格式

Python
@dataclass
class DiffResult:
    """Diff 查询结构化结果"""
    from_ref: str
    to_ref: str
    entity_type: str
    changes: list[DiffEntry]
    summary: DiffSummary
    metadata: DiffMetadata

    def to_table(self) -> pa.Table:
        """转换为 Arrow Table"""
        return pa.table({
            'entity_id': [c.entity_id for c in self.changes],
            'change_type': [c.change_type for c in self.changes],
            'old_value': [json.dumps(c.old_value) for c in self.changes],
            'new_value': [json.dumps(c.new_value) for c in self.changes],
            'changed_properties': [
                c.changed_properties for c in self.changes
            ],
        })

    def to_json(self) -> dict:
        """转换为 JSON"""
        return {
            'from': self.from_ref,
            'to': self.to_ref,
            'entity_type': self.entity_type,
            'summary': {
                'added': self.summary.added,
                'modified': self.summary.modified,
                'deleted': self.summary.deleted,
                'total': self.summary.total,
            },
            'changes': [
                {
                    'entity_id': c.entity_id,
                    'change_type': c.change_type,
                    'old_value': c.old_value,
                    'new_value': c.new_value,
                }
                for c in self.changes
            ]
        }


@dataclass
class DiffSummary:
    """Diff 摘要统计"""
    added: int
    modified: int
    deleted: int

    @property
    def total(self) -> int:
        return self.added + self.modified + self.deleted


@dataclass
class DiffMetadata:
    """Diff 元数据"""
    algorithm_used: str
    execution_time_ms: int
    files_scanned: int
    rows_compared: int
    from_snapshot_id: int
    to_snapshot_id: int

#5.2 Diff 可视化

Code
Diff 可视化输出示例:

DIFF Customer FROM '2024-03-01' TO '2024-03-31'

变更摘要:
  ██████████████████░░ 85% 未变更 (850/1000)
  ████░░░░░░░░░░░░░░░░ 8%  修改   (80)
  ██░░░░░░░░░░░░░░░░░░ 5%  新增   (50)
  █░░░░░░░░░░░░░░░░░░░ 2%  删除   (20)

属性变更热力图:
┌─────────────┬──────────┬──────────┐
│ 属性          │ 变更次数    │ 热度      │
├─────────────┼──────────┼──────────┤
│ credit_rating│ 45       │ ████████ │
│ risk_level   │ 30       │ █████    │
│ address      │ 25       │ ████     │
│ phone        │ 15       │ ██       │
│ name         │ 5        │ █        │
└─────────────┴──────────┴──────────┘

时间分布:
  3/01-3/07: ███████  35 变更
  3/08-3/14: █████    25 变更
  3/15-3/21: ████████ 40 变更
  3/22-3/28: █████████ 45 变更
  3/29-3/31: ██       5 变更

#6. 变更订阅

#6.1 基于 Diff 的变更通知

Python
class DiffBasedSubscription:
    """基于 Diff 的变更订阅系统"""

    def __init__(self, diff_engine: DiffEngine, notification_service: NotificationService):
        self._diff_engine = diff_engine
        self._notifier = notification_service
        self._subscriptions: list[Subscription] = []

    async def check_and_notify(self):
        """定期检查变更并通知订阅者"""
        for sub in self._subscriptions:
            last_check = sub.last_check_snapshot
            current = await self._get_current_snapshot(sub.entity_type)

            if current == last_check:
                continue

            diff = await self._diff_engine.diff(
                entity_type=sub.entity_type,
                from_snapshot=last_check,
                to_snapshot=current,
                columns=sub.watched_properties
            )

            if diff.summary.total > 0:
                matched = self._filter_changes(diff.changes, sub.filter)
                if matched:
                    await self._notifier.send(
                        subscriber=sub.subscriber,
                        notification=ChangeNotification(
                            entity_type=sub.entity_type,
                            changes=matched,
                            summary=DiffSummary(
                                added=sum(1 for c in matched if c.change_type == 'ADDED'),
                                modified=sum(1 for c in matched if c.change_type == 'MODIFIED'),
                                deleted=sum(1 for c in matched if c.change_type == 'DELETED'),
                            )
                        )
                    )

            sub.last_check_snapshot = current

#7. 性能优化

#7.1 Diff 缓存

Code
Diff 缓存策略:

问题:频繁的 Diff 查询重复读取相同快照

解决:分层缓存
  L1: 快照元数据缓存(Manifest List / File)
      → 避免重复读取 S3 元数据
      → TTL: 1 小时

  L2: 文件内容缓存(热数据文件的 Arrow 表示)
      → 避免重复反序列化 Parquet
      → LRU: 512 MB

  L3: Diff 结果缓存(相同快照对的 Diff 结果)
      → 完全避免重复计算
      → Key: (from_snap, to_snap, entity_type, columns)
      → TTL: 快照过期前永久有效

#7.2 并行 Diff

Python
class ParallelDiffExecutor:
    """并行 Diff 执行器"""

    async def diff_parallel(
        self,
        entity_type: str,
        from_snap: int,
        to_snap: int,
        parallelism: int = 8
    ) -> DiffResult:
        # 将数据文件分成 N 个分片
        file_groups = self._partition_files(
            entity_type, from_snap, to_snap, parallelism
        )

        # 并行执行 Diff
        tasks = [
            self._diff_file_group(group, from_snap, to_snap)
            for group in file_groups
        ]
        partial_results = await asyncio.gather(*tasks)

        # 合并结果
        return self._merge_diff_results(partial_results)

#8. 测试策略

Python
class TestDiffQueries:

    async def test_time_diff_detects_modifications(self):
        t1 = await snapshot_after(
            insert('Person', 'p1', {'name': 'Alice', 'age': 30})
        )
        t2 = await snapshot_after(
            update('p1', {'age': 31})
        )
        diff = await execute_diff('Person', t1, t2)
        assert diff.summary.modified == 1
        assert diff.changes[0].old_value['age'] == 30
        assert diff.changes[0].new_value['age'] == 31

    async def test_branch_diff(self):
        await create_branch('test-branch', from_ref='main')
        await on_branch('test-branch', lambda:
            insert('Person', 'p2', {'name': 'Bob'})
        )
        diff = await execute_diff_branch(
            'Person', 'main', 'test-branch'
        )
        assert diff.summary.added == 1

    async def test_three_way_merge_conflict(self):
        base = await current_snapshot()
        await on_branch('a', lambda: update('p1', {'age': 35}))
        await on_branch('b', lambda: update('p1', {'age': 40}))
        result = three_way_diff(base, 'a', 'b')
        assert result.has_conflicts
        assert 'age' in result.conflicts[0].conflicting_properties

    async def test_schema_evolution_diff(self):
        t1 = await current_snapshot()
        await add_column('Person', 'department', 'STRING')
        await update('p1', {'department': 'Engineering'})
        t2 = await current_snapshot()
        diff = await execute_diff('Person', t1, t2)
        assert any(
            sc.type == 'COLUMN_ADDED'
            for c in diff.changes for sc in c.schema_changes
        )

#Key Takeaways

  1. Diff 查询是"变更感知"的核心能力:不仅告诉你"数据是什么",还告诉你"数据变了什么",是审计合规、变更追踪和协作开发的基础。

  2. 三种 Diff 算法适配不同场景:全量对比用于小数据集,增量文件对比利用 Iceberg 快照差异提升大数据效率,Changelog 提供最精确的变更记录。

  3. Schema 演进下的 Diff 需要 field-id 匹配:使用 Iceberg 的 field-id 而非列名进行匹配,正确处理列重命名、新增、删除等 Schema 变更。

  4. 三方 Diff 支持分支合并冲突检测:基于共同祖先的三方对比可以精确识别属性级冲突,支持自动合并和人工解决。

  5. Diff 结果缓存是性能关键:相同快照对的 Diff 结果是不可变的,可以永久缓存直到快照过期。

#Next Article

下一篇 S3-12《分析引擎:OLAP 能力的 Ontology 封装》 将展示如何将 Doris 的 OLAP 能力封装为 Ontology 语义的分析查询,实现指标聚合、多维分析和实时仪表盘。

Tags: #DiffQuery #BranchDiff #ChangeTracking #Nessie #ThreeWayMerge #SchemaEvolution #ChangeSubscription #Audit #智策平台 #coomia-dip #数据基座