返回博客

联邦模式:跨组织的 Ontology 协作

企业集团中,不同业务线通常运行独立的平台实例。但跨组织的决策需要整合多源数据:

Coomia发布于 2025年12月21日11 分钟阅读
分享本文Twitter / X

联邦模式:跨组织的 Ontology 协作

系列:S10 设计模式 · 第 7 篇 | 难度:高级 | 阅读时间:18 分钟

#TL;DR

  • 联邦模式(Federation Pattern)允许多个独立的 coomia-dip 实例在保持数据主权的前提下进行协作查询和数据共享。
  • coomia-dip 实现了三级联邦:集群内联邦(跨 World)、组织内联邦(跨集群)、跨组织联邦(跨机构),每级都有不同的信任模型和访问控制。
  • 联邦查询通过查询分解(Query Decomposition)、远程执行和结果聚合实现,支持谓词下推优化以减少数据传输。

#引言:数据孤岛与协作需求

企业集团中,不同业务线通常运行独立的平台实例。但跨组织的决策需要整合多源数据:

Code
场景 1:集团风控需要整合子公司 A 的交易数据和子公司 B 的征信数据
场景 2:供应链优化需要联合上游供应商和下游分销商的库存数据
场景 3:跨境合规需要联合不同法域的 KYC/AML 数据,但数据不能出境

传统方案是 ETL 数据到中央仓库,但这面临数据主权、合规、实时性等问题。联邦模式提供了更好的选择:数据留在原地,查询走过去

#一、联邦架构

#1.1 三级联邦体系

Python
from dataclasses import dataclass, field
from enum import Enum
from typing import Any

class FederationLevel(Enum):
    """Federation levels with different trust models."""
    INTRA_CLUSTER = "intra_cluster"   # 集群内:跨 World 联邦
    INTRA_ORG = "intra_org"           # 组织内:跨集群联邦
    CROSS_ORG = "cross_org"           # 跨组织:跨机构联邦

@dataclass
class FederationNode:
    """A node in the federation network."""
    node_id: str
    name: str
    endpoint: str
    level: FederationLevel
    trust_level: str                   # full | limited | minimal
    capabilities: list[str] = field(default_factory=list)
    shared_object_types: list[str] = field(default_factory=list)
    access_policies: dict[str, Any] = field(default_factory=dict)
    metadata: dict[str, Any] = field(default_factory=dict)

@dataclass
class FederationRegistry:
    """Registry of federation nodes."""
    local_node_id: str
    nodes: dict[str, FederationNode] = field(default_factory=dict)

    def register(self, node: FederationNode) -> None:
        """Register a federation node."""
        self.nodes[node.node_id] = node

    def get_nodes_for_type(self, object_type: str) -> list[FederationNode]:
        """Find nodes that share a specific ObjectType."""
        return [
            n for n in self.nodes.values()
            if object_type in n.shared_object_types
        ]

#1.2 联邦目录服务

联邦目录记录每个节点共享的 ObjectType 和 Schema,类似 DNS 在互联网中的角色:

Python
class FederationCatalog:
    """Catalog service for federated Ontology schemas."""

    async def publish_schema(
        self,
        object_type: str,
        schema: dict,
        visibility: str = "federation",
    ) -> None:
        """Publish a schema to the federation catalog."""
        await self._catalog_store.save({
            "node_id": self._local_node_id,
            "object_type": object_type,
            "schema": schema,
            "visibility": visibility,
            "published_at": datetime.utcnow().isoformat(),
            "version": schema.get("version", 1),
        })

        # 通知联邦中的其他节点
        for node in self._registry.nodes.values():
            if node.node_id != self._local_node_id:
                await self._notify_schema_update(node, object_type, schema)

    async def discover_schema(
        self, object_type: str
    ) -> list[dict]:
        """Discover all nodes that provide a specific ObjectType."""
        results = []
        for node in self._registry.nodes.values():
            try:
                schema = await self._fetch_remote_schema(node, object_type)
                if schema:
                    results.append({
                        "node_id": node.node_id,
                        "node_name": node.name,
                        "schema": schema,
                        "trust_level": node.trust_level,
                    })
            except Exception:
                continue  # 节点不可用时跳过
        return results

    async def resolve_schema_conflicts(
        self, object_type: str, schemas: list[dict]
    ) -> dict:
        """Resolve schema conflicts across federation nodes."""
        if not schemas:
            raise ValueError(f"No schemas found for {object_type}")

        # 使用 Schema 合并策略
        base = schemas[0]["schema"]
        for s in schemas[1:]:
            base = self._merge_schemas(base, s["schema"])
        return base

    def _merge_schemas(self, schema_a: dict, schema_b: dict) -> dict:
        """Merge two schemas, taking the union of fields."""
        merged = dict(schema_a)
        for field_name, field_def in schema_b.get("properties", {}).items():
            if field_name not in merged.get("properties", {}):
                merged.setdefault("properties", {})[field_name] = field_def
        return merged

#二、联邦查询引擎

#2.1 查询分解

联邦查询引擎将一个查询分解为可以在各节点本地执行的子查询:

Python
@dataclass
class FederatedQueryPlan:
    """Execution plan for a federated query."""
    query_id: str
    original_query: str
    sub_queries: list["SubQuery"]
    aggregation: dict[str, Any]
    estimated_cost: float

@dataclass
class SubQuery:
    """A sub-query to be executed on a specific node."""
    node_id: str
    query: str
    object_type: str
    filters: list[dict]             # 推送到远程的过滤条件
    projections: list[str]          # 需要返回的字段
    estimated_rows: int

class FederatedQueryPlanner:
    """Plan federated queries across nodes."""

    async def plan(self, query: str) -> FederatedQueryPlan:
        """Create an execution plan for a federated query."""
        # 解析查询
        parsed = self._parser.parse(query)

        # 确定涉及的 ObjectType
        object_types = parsed.referenced_types

        # 为每个 ObjectType 找到数据源节点
        sub_queries = []
        for obj_type in object_types:
            nodes = self._registry.get_nodes_for_type(obj_type)
            for node in nodes:
                # 谓词下推优化
                pushable_filters = self._extract_pushable_filters(
                    parsed.filters, obj_type, node
                )

                sub_queries.append(SubQuery(
                    node_id=node.node_id,
                    query=self._build_sub_query(obj_type, pushable_filters, parsed.projections),
                    object_type=obj_type,
                    filters=pushable_filters,
                    projections=self._get_required_projections(parsed, obj_type),
                    estimated_rows=await self._estimate_rows(node, obj_type, pushable_filters),
                ))

        return FederatedQueryPlan(
            query_id=generate_id(),
            original_query=query,
            sub_queries=sub_queries,
            aggregation=self._plan_aggregation(parsed),
            estimated_cost=sum(sq.estimated_rows for sq in sub_queries),
        )

    def _extract_pushable_filters(
        self, filters: list, obj_type: str, node: FederationNode
    ) -> list:
        """Extract filters that can be pushed to a remote node."""
        pushable = []
        for f in filters:
            if f["field"].startswith(f"{obj_type}."):
                pushable.append(f)
        return pushable

#2.2 联邦查询执行器

Python
class FederatedQueryExecutor:
    """Execute federated queries across nodes."""

    async def execute(self, plan: FederatedQueryPlan) -> list[dict]:
        """Execute a federated query plan."""
        import asyncio

        # 并行执行所有子查询
        tasks = [
            self._execute_sub_query(sq) for sq in plan.sub_queries
        ]
        results = await asyncio.gather(*tasks, return_exceptions=True)

        # 处理失败的子查询
        valid_results = []
        for i, result in enumerate(results):
            if isinstance(result, Exception):
                node_id = plan.sub_queries[i].node_id
                # 记录失败但继续处理其他结果
                await self._log_sub_query_failure(node_id, result)
            else:
                valid_results.extend(result)

        # 聚合结果
        return self._aggregate_results(valid_results, plan.aggregation)

    async def _execute_sub_query(self, sub_query: SubQuery) -> list[dict]:
        """Execute a sub-query on a remote node."""
        node = self._registry.nodes[sub_query.node_id]

        if node.node_id == self._registry.local_node_id:
            # 本地执行
            return await self._local_executor.execute(sub_query.query)

        # 远程执行(通过 gRPC)
        async with grpc_channel(node.endpoint) as channel:
            stub = FederatedQueryServiceStub(channel)
            response = await stub.ExecuteQuery(
                FederatedQueryRequest(
                    query=sub_query.query,
                    object_type=sub_query.object_type,
                    filters=sub_query.filters,
                    projections=sub_query.projections,
                    caller_node_id=self._registry.local_node_id,
                )
            )
            return [dict(row) for row in response.rows]

    def _aggregate_results(
        self, results: list[dict], aggregation: dict
    ) -> list[dict]:
        """Aggregate results from multiple sub-queries."""
        if aggregation.get("type") == "union":
            return results
        elif aggregation.get("type") == "join":
            return self._join_results(
                results, aggregation["join_key"]
            )
        elif aggregation.get("type") == "group_by":
            return self._group_results(
                results, aggregation["group_key"], aggregation["agg_func"]
            )
        return results

#三、数据主权与访问控制

#3.1 联邦访问策略

每个节点可以精确控制哪些数据可以被联邦查询访问:

Python
@dataclass
class FederationAccessPolicy:
    """Access policy for federated data sharing."""
    policy_id: str
    object_type: str
    allowed_nodes: list[str] | None = None    # None = 所有节点
    denied_nodes: list[str] = field(default_factory=list)
    allowed_fields: list[str] | None = None   # None = 所有字段
    denied_fields: list[str] = field(default_factory=list)
    row_filter: str | None = None             # 行级过滤表达式
    data_masking: dict[str, str] = field(default_factory=dict)  # 字段脱敏规则
    max_rows: int | None = None               # 最大返回行数
    audit_required: bool = True

class FederationAccessController:
    """Enforce federation access policies."""

    async def check_access(
        self,
        caller_node_id: str,
        object_type: str,
        requested_fields: list[str],
    ) -> dict:
        """Check if a federated query is allowed."""
        policy = await self._get_policy(object_type)

        if not policy:
            return {"allowed": False, "reason": "No federation policy defined"}

        # 检查节点访问权限
        if policy.denied_nodes and caller_node_id in policy.denied_nodes:
            return {"allowed": False, "reason": "Node is denied"}

        if policy.allowed_nodes and caller_node_id not in policy.allowed_nodes:
            return {"allowed": False, "reason": "Node is not in allowed list"}

        # 检查字段访问权限
        denied_fields = set(requested_fields) & set(policy.denied_fields)
        if denied_fields:
            return {
                "allowed": False,
                "reason": f"Denied fields: {denied_fields}",
            }

        # 应用数据脱敏
        masking_rules = {}
        for field_name in requested_fields:
            if field_name in policy.data_masking:
                masking_rules[field_name] = policy.data_masking[field_name]

        return {
            "allowed": True,
            "masking_rules": masking_rules,
            "row_filter": policy.row_filter,
            "max_rows": policy.max_rows,
        }

    async def apply_masking(
        self, data: list[dict], masking_rules: dict[str, str]
    ) -> list[dict]:
        """Apply data masking to query results."""
        masked = []
        for row in data:
            masked_row = dict(row)
            for field_name, mask_type in masking_rules.items():
                if field_name in masked_row:
                    masked_row[field_name] = self._mask_value(
                        masked_row[field_name], mask_type
                    )
            masked.append(masked_row)
        return masked

    def _mask_value(self, value: Any, mask_type: str) -> Any:
        """Mask a value according to the masking rule."""
        if mask_type == "hash":
            import hashlib
            return hashlib.sha256(str(value).encode()).hexdigest()[:16]
        elif mask_type == "partial":
            s = str(value)
            return s[:2] + "*" * (len(s) - 4) + s[-2:] if len(s) > 4 else "****"
        elif mask_type == "null":
            return None
        elif mask_type == "category":
            # 保留类别信息但隐藏具体值
            return f"[MASKED:{type(value).__name__}]"
        return value

#3.2 跨境数据合规

对于跨境联邦查询,coomia-dip 确保数据不出境,只传输聚合结果:

Python
class CrossBorderFederationPolicy:
    """Ensure compliance with cross-border data regulations."""

    async def enforce(
        self,
        query_plan: FederatedQueryPlan,
        data_residency_rules: dict[str, str],
    ) -> FederatedQueryPlan:
        """Modify query plan to comply with data residency rules."""
        modified_sub_queries = []

        for sq in query_plan.sub_queries:
            node = self._registry.nodes[sq.node_id]
            node_region = node.metadata.get("region", "unknown")
            caller_region = self._local_region

            if node_region != caller_region:
                # 跨区域查询——只允许聚合结果
                if not self._is_aggregation_only(sq):
                    # 修改子查询为聚合查询
                    sq = self._convert_to_aggregation(sq)

                # 添加数据脱敏
                sq.filters.append({
                    "type": "masking",
                    "policy": "cross_border_default",
                })

            modified_sub_queries.append(sq)

        query_plan.sub_queries = modified_sub_queries
        return query_plan

#四、联邦数据同步

#4.1 事件驱动的联邦同步

对于需要近实时同步的场景,coomia-dip 使用事件驱动的联邦同步:

Python
class FederationSyncService:
    """Event-driven synchronization across federation nodes."""

    async def publish_change(
        self,
        object_type: str,
        change_type: str,
        data: dict,
    ) -> None:
        """Publish a change event to interested federation nodes."""
        interested_nodes = self._registry.get_nodes_for_type(object_type)

        for node in interested_nodes:
            if node.node_id == self._registry.local_node_id:
                continue

            # 应用访问策略过滤
            access = await self._access_controller.check_access(
                node.node_id, object_type, list(data.keys())
            )
            if not access["allowed"]:
                continue

            # 脱敏后发送
            if access.get("masking_rules"):
                data = (await self._access_controller.apply_masking(
                    [data], access["masking_rules"]
                ))[0]

            await self._event_bus.publish(
                topic=f"federation.{node.node_id}.{object_type}",
                event={
                    "source_node": self._registry.local_node_id,
                    "object_type": object_type,
                    "change_type": change_type,
                    "data": data,
                    "timestamp": datetime.utcnow().isoformat(),
                },
            )

    async def subscribe_changes(
        self,
        source_node_id: str,
        object_type: str,
        handler: callable,
    ) -> None:
        """Subscribe to changes from a federation node."""
        topic = f"federation.{self._registry.local_node_id}.{object_type}"
        await self._event_bus.subscribe(topic, handler)

#五、实战场景

#5.1 集团风控联邦查询

Python
# 查询跨子公司的风控数据
federated_query = """
SELECT
    t.customer_id,
    t.transaction_amount,
    c.credit_score,
    r.risk_level
FROM
    SubsidiaryA.Transaction t
    JOIN SubsidiaryB.CreditRecord c ON t.customer_id = c.customer_id
    JOIN Headquarters.RiskAssessment r ON t.customer_id = r.customer_id
WHERE
    t.transaction_amount > 100000
    AND c.credit_score < 600
"""

# 联邦查询引擎自动分解为三个子查询分别发送到三个节点

#5.2 供应链联邦优化

Python
# 跨组织的供应链库存联邦查询
supply_chain_query = """
SELECT
    s.supplier_id,
    s.product_id,
    s.available_qty,
    d.demand_forecast,
    d.demand_forecast - s.available_qty AS gap
FROM
    Supplier.Inventory s
    JOIN Internal.DemandForecast d ON s.product_id = d.product_id
WHERE
    d.demand_forecast > s.available_qty
ORDER BY gap DESC
"""

#Key Takeaways

  1. 数据主权:联邦模式保证数据留在原处,只有查询和聚合结果跨网络传输
  2. 三级联邦:集群内、组织内、跨组织三级联邦,信任模型逐级递减
  3. 查询优化:谓词下推和并行执行减少数据传输和查询延迟
  4. 精细访问控制:字段级、行级、节点级的多维访问控制和数据脱敏
  5. 跨境合规:自动确保跨境查询只传输聚合结果,原始数据不出境
  6. 事件同步:近实时的事件驱动联邦同步,支持增量更新

#Next Article

下一篇我们将探讨级联模式(Cascade Pattern)——coomia-dip 如何处理 Ontology 对象之间的级联更新、级联删除和依赖传播。

S10-08: 级联模式:依赖传播与影响分析

#Tags

#设计模式 #联邦模式 #Federation #跨组织 #数据主权 #联邦查询 #访问控制 #数据脱敏 #跨境合规