返回博客

World Transform:全局数据一致性变换

Tags: #WorldTransform #Consistency #GlobalState #Transaction #Ontology #智策平台

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

系列: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

Code
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 模型

Code
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

Python
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 依赖图

Code
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 原子性保障

Python
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 冲突检测

Code
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

Python
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 增量检测

Code
增量检测机制:

基于 Iceberg 快照差异:
  上次执行快照 ID: snap-1001
  当前快照 ID: snap-1005

  变更检测:
    snap-1001 → snap-1005 之间新增/修改/删除的数据文件
    → 提取变更的 entity_id 集合
    → 只处理这些 entity_id

效果:
  全量执行:处理 100,000 客户 → 5 分钟
  增量执行:处理 500 变更客户 → 3 秒
  加速比:100 倍

#5. 输出验证

#5.1 验证规则

Python
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. 性能优化

Code
World Transform 性能优化策略:

1. 并行分片执行
   将 Entity 按 ID 范围分片,各分片并行执行
   10 个分片 × 10,000 实体 = 10 倍加速

2. 向量化计算
   使用 Arrow/NumPy 进行批量计算
   逐行计算 → 批量向量计算 = 5-10 倍加速

3. 缓存读取
   频繁访问的实体和关系缓存到内存
   避免重复 Doris 查询

4. 增量执行(前文已述)
   只处理变更数据 = 10-100 倍加速

综合效果:
  全量非优化:30 分钟
  全量优化后:3 分钟
  增量优化后:3 秒

#7. 测试策略

Python
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

  1. World Transform 提供了全局一致性的数据变换能力:单次事务中同时修改多个 Entity Type 和关系,通过 Nessie 分支实现原子性——要么全部成功,要么全部回滚。

  2. 声明式 Input/Output 自动构建依赖图:Transform 只需声明读取和写入的 Entity Type,系统自动推导执行顺序、检测循环依赖。

  3. 增量执行是性能的关键:通过 Iceberg 快照差异检测变更数据,只处理受影响的实体,将执行时间从分钟级降到秒级。

  4. 输出验证防止了误操作:Schema 兼容性、参照完整性、业务规则和变更量阈值检查,确保 Transform 输出的正确性。

  5. Nessie 分支提供了天然的事务隔离:每次 Transform 在独立分支上执行,合并时检测冲突,支持属性级别的自动合并。

#Next Article

下一篇 S3-22《订阅系统:实时数据变更通知》 将展示如何基于 Ontology 变更事件构建实时订阅和通知系统。

Tags: #WorldTransform #Consistency #GlobalState #AtomicTransaction #NessieBranch #IncrementalExecution #DependencyGraph #智策平台 #coomia-dip #数据基座