World Transform:全局数据一致性变换
Tags: #WorldTransform #Consistency #GlobalState #Transaction #Ontology #智策平台
“系列:S3 数据基座 · 第 21 篇 | 难度:高级 | 阅读时间:20 分钟
World Transform:全局数据一致性变换
Tags: #WorldTransform #Consistency #GlobalState #Transaction #Ontology #智策平台
#TL;DR
World Transform 是 Palantir Foundry 的核心概念之一——对整个 Ontology 世界状态进行原子性的全局变换。不同于单个 Pipeline 只处理一种数据源,World Transform 可以在单次事务中同时修改多个 Entity Type 的数据、更新关系图、重新计算指标,并保证全局一致性。本文完整解析 World Transform 的概念模型、事务管理机制、依赖图计算、增量执行策略、Nessie 分支集成以及与 Pipeline DSL 的协作方式。
#1. World Transform 概念
#1.1 什么是 World Transform
World Transform vs Pipeline Transform:
Pipeline Transform(局部变换):
输入:一个数据源(如 MySQL 表)
输出:一个目标(如 Ontology Entity Type)
范围:单一数据流
事务:Pipeline 内部一致
World Transform(全局变换):
输入:整个 Ontology 世界状态(多个 Entity Type + 关系图)
输出:新的世界状态(原子更新)
范围:跨 Entity Type、跨关系的全局计算
事务:全局一致(要么全部成功,要么全部回滚)
示例:
"重新计算所有客户的风险等级"
→ 需要读取 Customer、Transaction、RiskModel
→ 更新 Customer.risk_level
→ 更新 RiskAlert 关系
→ 更新 RiskScore 指标
→ 以上必须原子完成
#1.2 Palantir Foundry 的 Transform 模型
Foundry Transform 模型:
┌──────────────────────────────────────────┐
│ World State v1 │
│ ┌────────┐ ┌────────┐ ┌────────┐ │
│ │Customer│ │Product │ │Order │ │
│ │ (1000) │ │ (200) │ │ (5000) │ │
│ └────────┘ └────────┘ └────────┘ │
│ ┌────────────────────────────┐ │
│ │ Relationships (10000) │ │
│ └────────────────────────────┘ │
└──────────────────┬───────────────────────┘
│ World Transform
│ (atomic)
▼
┌──────────────────────────────────────────┐
│ World State v2 │
│ ┌────────┐ ┌────────┐ ┌────────┐ │
│ │Customer│ │Product │ │Order │ │
│ │ (1000) │ │ (200) │ │ (5200) │ │
│ └────────┘ └────────┘ └────────┘ │
│ ┌────────────────────────────┐ │
│ │ Relationships (10500) │ │
│ └────────────────────────────┘ │
│ ┌────────────────────────────┐ │
│ │ Derived: RiskScores (1000)│ │
│ └────────────────────────────┘ │
└──────────────────────────────────────────┘
#2. World Transform 定义
#2.1 声明式 API
from onto_transform import WorldTransform, Input, Output
@WorldTransform(
name="risk-scoring",
description="重新计算所有客户的风险评分",
schedule="0 2 * * *", # 每天凌晨 2 点
version="2.0.0",
)
class RiskScoringTransform:
"""风险评分全局变换"""
# 声明输入依赖
customers = Input("Customer", properties=["id", "credit_rating", "region"])
transactions = Input("Transaction", properties=["customer_id", "amount", "type"])
risk_model = Input("RiskModel", properties=["id", "weights", "thresholds"])
# 声明输出
risk_scores = Output("Customer", properties=["risk_level", "risk_score"])
risk_alerts = Output("RiskAlert", relationship=True)
def compute(self, ctx: TransformContext):
"""执行变换逻辑"""
# 读取输入
customers = ctx.read(self.customers)
transactions = ctx.read(self.transactions)
model = ctx.read(self.risk_model).first()
# 计算每个客户的风险评分
for customer in customers:
customer_txns = transactions.filter(
customer_id=customer.id
)
risk_score = self._calculate_risk(
customer, customer_txns, model
)
# 更新客户风险等级
ctx.update(self.risk_scores, customer.id, {
'risk_level': self._score_to_level(risk_score),
'risk_score': risk_score,
})
# 如果风险高,创建告警关系
if risk_score > model.thresholds['high']:
ctx.create_relationship(self.risk_alerts, {
'source_id': customer.id,
'target_id': f"alert-{customer.id}-{ctx.run_id}",
'alert_type': 'HIGH_RISK',
'score': risk_score,
})
def _calculate_risk(self, customer, transactions, model):
weights = model.weights
total_amount = sum(t.amount for t in transactions)
txn_count = len(transactions)
high_value_count = sum(
1 for t in transactions if t.amount > 10000
)
score = (
weights['amount'] * min(total_amount / 100000, 1.0) +
weights['frequency'] * min(txn_count / 100, 1.0) +
weights['high_value'] * min(high_value_count / 10, 1.0) +
weights['credit'] * (1.0 - self._credit_to_score(customer.credit_rating))
)
return round(score * 100, 2)
#2.2 依赖图
World Transform 依赖图:
Transform 定义了读取(Input)和写入(Output)的 Entity Type
系统自动构建依赖图:
risk-model-training
Output: RiskModel
│
▼
risk-scoring (本 Transform)
Input: Customer, Transaction, RiskModel
Output: Customer.risk_level, RiskAlert
│
▼
risk-reporting
Input: Customer.risk_level, RiskAlert
Output: RiskReport
customer-segmentation
Input: Customer.risk_level
Output: Customer.segment
依赖图确保:
1. risk-model-training 先于 risk-scoring 执行
2. risk-scoring 先于 risk-reporting 和 customer-segmentation
3. 循环依赖被检测并报错
#3. 事务管理
#3.1 原子性保障
class WorldTransformExecutor:
"""World Transform 事务执行器"""
async def execute(
self, transform: WorldTransform
) -> TransformResult:
# 步骤 1:创建 Nessie 分支
branch_name = f"wt-{transform.name}-{uuid4().hex[:8]}"
await self._nessie.create_branch(
branch_name, from_ref="main"
)
try:
# 步骤 2:在分支上执行变换
ctx = TransformContext(
branch=branch_name,
nessie=self._nessie,
iceberg=self._iceberg,
)
transform.compute(ctx)
# 步骤 3:验证输出
validation = await self._validate_output(ctx)
if not validation.passed:
raise TransformValidationError(validation.errors)
# 步骤 4:合并到 main(原子操作)
await self._nessie.merge(
from_branch=branch_name,
to_branch="main",
conflict_resolution="REJECT", # 有冲突则失败
)
return TransformResult(
status="SUCCESS",
branch=branch_name,
changes=ctx.get_change_summary(),
)
except Exception as e:
# 回滚:删除分支(所有变更自动丢弃)
await self._nessie.delete_branch(branch_name)
return TransformResult(
status="FAILED",
error=str(e),
)
#3.2 冲突检测
World Transform 冲突场景:
场景:两个 Transform 并发修改同一个 Customer
main ──●──────────────────── HEAD
\
\── wt-risk-scoring (修改 Customer.risk_level)
\
\── wt-segmentation (修改 Customer.segment)
合并策略:
1. 先完成的 Transform 合并到 main ✓
2. 后完成的检测到冲突(Customer 已被修改)
3. 冲突处理选项:
a) REJECT:直接失败,让调度器重试
b) RETRY:重新读取最新数据,重新执行
c) PROPERTY_LEVEL:不同属性不冲突
risk_level 和 segment 是不同属性 → 可以自动合并
#4. 增量执行
#4.1 增量 Transform
class IncrementalWorldTransform:
"""增量 World Transform — 只处理变更数据"""
@WorldTransform(
name="risk-scoring-incremental",
incremental=True,
)
class IncrementalRiskScoring:
customers = Input("Customer", incremental=True)
transactions = Input("Transaction", incremental=True)
def compute(self, ctx: TransformContext):
# 只获取自上次执行以来的变更
changed_customers = ctx.read_changes(self.customers)
new_transactions = ctx.read_changes(self.transactions)
# 确定需要重新计算的客户集合
affected_ids = set()
affected_ids.update(c.id for c in changed_customers)
affected_ids.update(t.customer_id for t in new_transactions)
# 只重新计算受影响的客户
for cid in affected_ids:
customer = ctx.read_entity("Customer", cid)
transactions = ctx.read_related(
"Transaction", customer_id=cid
)
risk_score = self._calculate_risk(customer, transactions)
ctx.update(self.risk_scores, cid, {
'risk_level': self._score_to_level(risk_score),
'risk_score': risk_score,
})
#4.2 增量检测
增量检测机制:
基于 Iceberg 快照差异:
上次执行快照 ID: snap-1001
当前快照 ID: snap-1005
变更检测:
snap-1001 → snap-1005 之间新增/修改/删除的数据文件
→ 提取变更的 entity_id 集合
→ 只处理这些 entity_id
效果:
全量执行:处理 100,000 客户 → 5 分钟
增量执行:处理 500 变更客户 → 3 秒
加速比:100 倍
#5. 输出验证
#5.1 验证规则
class TransformOutputValidator:
"""Transform 输出验证器"""
def validate(self, ctx: TransformContext) -> ValidationResult:
errors = []
# 规则 1:输出 Schema 兼容性
for output in ctx.outputs:
schema_errors = self._check_schema(output)
errors.extend(schema_errors)
# 规则 2:参照完整性
for relationship in ctx.new_relationships:
if not ctx.entity_exists(relationship.source_id):
errors.append(f"Source entity {relationship.source_id} not found")
if not ctx.entity_exists(relationship.target_id):
errors.append(f"Target entity {relationship.target_id} not found")
# 规则 3:业务规则
for entity in ctx.updated_entities:
if entity.type == 'Customer':
if entity.properties.get('risk_score', 0) < 0:
errors.append(f"Negative risk score for {entity.id}")
if entity.properties.get('risk_score', 0) > 100:
errors.append(f"Risk score > 100 for {entity.id}")
# 规则 4:变更量检查(防止误操作)
change_ratio = ctx.change_count / ctx.total_entities
if change_ratio > 0.5:
errors.append(
f"Change ratio {change_ratio:.1%} exceeds safety threshold 50%. "
f"This may indicate a bug in the transform."
)
return ValidationResult(
passed=len(errors) == 0,
errors=errors
)
#6. 性能优化
World Transform 性能优化策略:
1. 并行分片执行
将 Entity 按 ID 范围分片,各分片并行执行
10 个分片 × 10,000 实体 = 10 倍加速
2. 向量化计算
使用 Arrow/NumPy 进行批量计算
逐行计算 → 批量向量计算 = 5-10 倍加速
3. 缓存读取
频繁访问的实体和关系缓存到内存
避免重复 Doris 查询
4. 增量执行(前文已述)
只处理变更数据 = 10-100 倍加速
综合效果:
全量非优化:30 分钟
全量优化后:3 分钟
增量优化后:3 秒
#7. 测试策略
class TestWorldTransform:
async def test_atomic_rollback_on_failure(self):
"""失败时所有变更应该回滚"""
initial_state = await get_world_state()
with pytest.raises(TransformError):
await execute_transform(FailingTransform())
final_state = await get_world_state()
assert initial_state == final_state # 无变更
async def test_incremental_only_processes_changes(self):
# 初始全量执行
await execute_transform(RiskScoringIncremental())
# 修改 5 个客户
for i in range(5):
await update_entity(f"customer-{i}", {"credit_rating": "B"})
# 增量执行应只处理 5 个客户
result = await execute_transform(RiskScoringIncremental())
assert result.processed_count == 5
async def test_concurrent_transforms_conflict(self):
"""并发 Transform 应检测冲突"""
t1 = asyncio.create_task(
execute_transform(RiskScoring())
)
t2 = asyncio.create_task(
execute_transform(Segmentation())
)
results = await asyncio.gather(t1, t2, return_exceptions=True)
# 至少一个应成功
successes = [r for r in results if not isinstance(r, Exception)]
assert len(successes) >= 1
async def test_validation_catches_bad_output(self):
result = await execute_transform(BadRiskScoring())
assert result.status == "FAILED"
assert "risk_score" in str(result.error)
#Key Takeaways
-
World Transform 提供了全局一致性的数据变换能力:单次事务中同时修改多个 Entity Type 和关系,通过 Nessie 分支实现原子性——要么全部成功,要么全部回滚。
-
声明式 Input/Output 自动构建依赖图:Transform 只需声明读取和写入的 Entity Type,系统自动推导执行顺序、检测循环依赖。
-
增量执行是性能的关键:通过 Iceberg 快照差异检测变更数据,只处理受影响的实体,将执行时间从分钟级降到秒级。
-
输出验证防止了误操作:Schema 兼容性、参照完整性、业务规则和变更量阈值检查,确保 Transform 输出的正确性。
-
Nessie 分支提供了天然的事务隔离:每次 Transform 在独立分支上执行,合并时检测冲突,支持属性级别的自动合并。
#Next Article
下一篇 S3-22《订阅系统:实时数据变更通知》 将展示如何基于 Ontology 变更事件构建实时订阅和通知系统。
Tags: #WorldTransform #Consistency #GlobalState #AtomicTransaction #NessieBranch #IncrementalExecution #DependencyGraph #智策平台 #coomia-dip #数据基座