Diff 查询:分支对比与变更追踪
Tags: #DiffQuery #BranchDiff #ChangeTracking #Nessie #Audit #智策平台
“系列: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 查询类型
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 分支模型回顾
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 算法一:全量对比
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)
全量对比算法特性:
时间复杂度:O(N + M),N 和 M 为两个版本的行数
空间复杂度:O(N + M)
优点:简单、正确
缺点:需要读取全部数据,数据量大时慢
适用场景:
- 数据量 < 100K 行
- 首次 Diff(无历史增量信息)
#2.2 算法二:增量文件对比
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)
增量文件对比算法特性:
时间复杂度:O(C),C 为变更的行数
空间复杂度:O(C)
优点:只读取变更部分,大数据集高效
缺点:依赖 Iceberg 文件级追踪
适用场景:
- 数据量 > 100K 行
- 变更量 < 总量的 10%
- 同一个 Iceberg 表的两个快照对比
#2.3 算法三:Changelog 对比
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 变更类型
Schema 演进对 Diff 的影响:
┌─────────────────┬────────────────────────────────────┐
│ Schema 变更 │ Diff 处理策略 │
├─────────────────┼────────────────────────────────────┤
│ 新增列 │ 新版本有值、旧版本填充 NULL │
│ 删除列 │ 旧版本有值、新版本标记为 DROPPED │
│ 列重命名 │ 通过 Iceberg field-id 映射 │
│ 类型变更 │ 自动类型提升(int→long, float→double)│
│ 新增 Entity Type │ 全部标记为 ADDED │
│ 删除 Entity Type │ 全部标记为 DELETED │
└─────────────────┴────────────────────────────────────┘
#3.2 兼容 Diff 实现
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
三方 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 冲突检测实现
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 冲突解决策略
冲突解决策略:
┌──────────────────┬──────────────────────────────────┐
│ 策略 │ 描述 │
├──────────────────┼──────────────────────────────────┤
│ TAKE_SOURCE │ 以源分支(被合并方)为准 │
│ TAKE_TARGET │ 以目标分支(合并到)为准 │
│ TAKE_LATEST │ 取最后修改时间较新的版本 │
│ MANUAL │ 人工选择 │
│ PROPERTY_LEVEL │ 属性级别自动合并 │
│ │ (不同属性分别取各自修改方) │
│ CUSTOM_FUNCTION │ 自定义合并函数 │
└──────────────────┴──────────────────────────────────┘
#5. Diff 结果结构化输出
#5.1 输出格式
@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 可视化
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 的变更通知
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 缓存
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
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. 测试策略
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
-
Diff 查询是"变更感知"的核心能力:不仅告诉你"数据是什么",还告诉你"数据变了什么",是审计合规、变更追踪和协作开发的基础。
-
三种 Diff 算法适配不同场景:全量对比用于小数据集,增量文件对比利用 Iceberg 快照差异提升大数据效率,Changelog 提供最精确的变更记录。
-
Schema 演进下的 Diff 需要 field-id 匹配:使用 Iceberg 的 field-id 而非列名进行匹配,正确处理列重命名、新增、删除等 Schema 变更。
-
三方 Diff 支持分支合并冲突检测:基于共同祖先的三方对比可以精确识别属性级冲突,支持自动合并和人工解决。
-
Diff 结果缓存是性能关键:相同快照对的 Diff 结果是不可变的,可以永久缓存直到快照过期。
#Next Article
下一篇 S3-12《分析引擎:OLAP 能力的 Ontology 封装》 将展示如何将 Doris 的 OLAP 能力封装为 Ontology 语义的分析查询,实现指标聚合、多维分析和实时仪表盘。
Tags: #DiffQuery #BranchDiff #ChangeTracking #Nessie #ThreeWayMerge #SchemaEvolution #ChangeSubscription #Audit #智策平台 #coomia-dip #数据基座