返回博客

源码精读:SchemaRegistryService — 本体注册的状态机

SchemaRegistryService 是 Control Layer(Control Layer)中管理本体 Schema 全生命周期的核心服务。它基于 Spring Boot 3.x + gRPC 构建,实现了 DRAFT → ACTIVE → DEPRECATED → ARCHIVED 的四态状态机,内置向前/向后兼容性检查、依赖分析、版本历史管理和批量导入导出。本文将深入分析状态机的转换规则、CompatibilityChecker 的三级兼容策略、SchemaDependencyChecker 的级联影响分析,以及 SchemaVersionHistory 的不可变版本链设计。

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

源码精读:SchemaRegistryService — 本体注册的状态机

系列:S9 源码精读 · 第 2 篇 | 难度:高级 | 阅读时间:25 分钟

#TL;DR

SchemaRegistryService 是 Control Layer(Control Layer)中管理本体 Schema 全生命周期的核心服务。它基于 Spring Boot 3.x + gRPC 构建,实现了 DRAFT → ACTIVE → DEPRECATED → ARCHIVED 的四态状态机,内置向前/向后兼容性检查、依赖分析、版本历史管理和批量导入导出。本文将深入分析状态机的转换规则、CompatibilityChecker 的三级兼容策略、SchemaDependencyChecker 的级联影响分析,以及 SchemaVersionHistory 的不可变版本链设计。

#目录

  1. 架构定位与模块结构
  2. 领域模型:SchemaEntity 与 SchemaStatus
  3. gRPC 服务入口:SchemaRegistryServiceImpl
  4. 状态机:四态转换与守卫条件
  5. 兼容性检查:CompatibilityChecker 三级策略
  6. 依赖分析:SchemaDependencyChecker
  7. 版本管理:SchemaVersionHistory
  8. Schema CRUD 完整流程
  9. 批量导入导出:YAML/JSON 双格式
  10. EventRef 验证与跨实体引用
  11. Key Takeaways

#1. 架构定位与模块结构

Schema Registry 位于 Control Layer,是整个平台的"元数据大脑"。所有 Layer 在操作数据之前,都需要先向 Schema Registry 查询实体类型定义。

Code
control-Layer/src/main/java/com/onto/control/schema/
├── api/
│   ├── SchemaRegistryServiceImpl.java    # gRPC 服务实现
│   ├── ActionTypeGrpcService.java        # ActionType 注册
│   └── StructTypeGrpcService.java        # StructType 注册
├── domain/
│   ├── SchemaEntity.java                 # Schema 聚合根
│   ├── SchemaStatus.java                 # 状态枚举
│   ├── SchemaLifecycleEvent.java         # 生命周期事件
│   └── SchemaVersionHistory.java         # 版本历史
├── service/
│   ├── SchemaService.java                # 业务服务
│   ├── CompatibilityChecker.java         # 兼容性检查器
│   ├── SchemaDependencyChecker.java      # 依赖检查器
│   └── EventRefValidator.java            # 事件引用验证
├── mapper/
│   └── SchemaMapper.java                 # Proto <-> Domain 映射
├── analyzer/
│   ├── DependencyAnalyzer.java           # 依赖图分析
│   ├── DependencyGraph.java              # 图数据结构
│   ├── ImpactAnalysis.java               # 影响评估
│   └── DependencyCacheService.java       # 依赖缓存
└── repository/
    └── SchemaRepository.java             # JPA 持久层

与 Palantir Foundry 的对标:Palantir Foundry 的 Ontology Manager 使用 REST API 管理 ObjectType、LinkType、ActionType。coomia-dip 的 SchemaRegistryService 对标此功能,但增加了状态机、兼容性检查和依赖分析——这些在 Foundry 中需要手动管理。

#2. 领域模型:SchemaEntity 与 SchemaStatus

Java
@Entity
@Table(name = "schema_registry")
public class SchemaEntity {
    @Id
    private String schemaId;
    private String schemaName;
    private String entityType;         // OBJECT_TYPE, LINK_TYPE, INTERFACE
    private int version;
    private String definition;          // JSON/YAML 格式的 Schema 定义
    private String worldId;

    @Enumerated(EnumType.STRING)
    private SchemaStatus status;

    private String createdBy;
    private LocalDateTime createdAt;
    private String updatedBy;
    private LocalDateTime updatedAt;

    // 兼容性配置
    @Enumerated(EnumType.STRING)
    private CompatibilityMode compatibilityMode;  // NONE, BACKWARD, FORWARD, FULL

    // 版本历史(JSON 存储)
    @Column(columnDefinition = "TEXT")
    private String versionHistory;
}
Java
public enum SchemaStatus {
    DRAFT,        // 草稿,可自由修改
    ACTIVE,       // 已激活,修改需要兼容性检查
    DEPRECATED,   // 已弃用,只读
    ARCHIVED      // 已归档,不可见
}

状态转换图

Code
DRAFT ──activate()──→ ACTIVE ──deprecate()──→ DEPRECATED ──archive()──→ ARCHIVED
  ↑                     │                        │
  │                     │ revert()               │ revert()
  └─────────────────────┘                        │
  ↑                                              │
  └──────────────────────────────────────────────┘

设计考量:DRAFT 状态允许自由修改而不触发兼容性检查。只有从 DRAFT → ACTIVE 的 activate() 转换会执行完整的验证链。这使得 Schema 的迭代开发更高效——在 DRAFT 阶段可以反复调整字段。

#3. gRPC 服务入口:SchemaRegistryServiceImpl

Java
@GrpcService
@Slf4j
@RequiredArgsConstructor
public class SchemaRegistryServiceImpl
        extends SchemaRegistryServiceGrpc.SchemaRegistryServiceImplBase {

    private final SchemaService schemaService;
    private final CompatibilityChecker compatibilityChecker;
    private final SchemaDependencyChecker dependencyChecker;
    private final EventRefValidator eventRefValidator;
    private final SchemaMapper mapper;

    @Override
    public void registerSchema(RegisterSchemaRequest request,
            StreamObserver<SchemaResponse> responseObserver) {
        try {
            String worldId = WorldContextHolder.getWorldId();
            String userId = WorldContextHolder.getUserId();

            log.info("RegisterSchema: name={}, world={}",
                request.getSchema().getSchemaName(), worldId);

            // 1. 转换 Proto → Domain
            SchemaEntity entity = mapper.toDomain(
                request.getSchema(), worldId, userId);

            // 2. 初始状态设置
            entity.setStatus(SchemaStatus.DRAFT);
            entity.setVersion(1);

            // 3. EventRef 验证(如果引用了事件类型)
            if (request.getSchema().hasEventRef()) {
                eventRefValidator.validate(
                    request.getSchema().getEventRef(), worldId);
            }

            // 4. 持久化
            SchemaEntity saved = schemaService.register(entity);

            // 5. 响应
            responseObserver.onNext(
                mapper.toResponse(saved));
            responseObserver.onCompleted();
        } catch (Exception e) {
            handleException(e, responseObserver);
        }
    }
}

注意:新注册的 Schema 始终以 DRAFT 状态开始,版本号从 1 开始。这确保了没有 Schema 可以跳过审查直接进入生产。

#4. 状态机:四态转换与守卫条件

Java
@Override
public void activateSchema(ActivateSchemaRequest request,
        StreamObserver<SchemaResponse> responseObserver) {
    try {
        String worldId = WorldContextHolder.getWorldId();
        SchemaEntity entity = schemaService.getById(
            request.getSchemaId(), worldId);

        // 守卫条件 1:只有 DRAFT 可以被激活
        if (entity.getStatus() != SchemaStatus.DRAFT) {
            throw Status.FAILED_PRECONDITION
                .withDescription(String.format(
                    "Schema '%s' is in %s state, only DRAFT can be activated",
                    entity.getSchemaName(), entity.getStatus()))
                .asRuntimeException();
        }

        // 守卫条件 2:兼容性检查(如果存在前一个 ACTIVE 版本)
        Optional<SchemaEntity> previousActive =
            schemaService.findActiveByName(
                entity.getSchemaName(), worldId);
        if (previousActive.isPresent()) {
            CompatibilityResult result =
                compatibilityChecker.check(
                    previousActive.get(), entity);
            if (!result.isCompatible()) {
                throw Status.FAILED_PRECONDITION
                    .withDescription(
                        "Compatibility check failed: "
                        + result.getViolations())
                    .asRuntimeException();
            }
            // 将旧版本转为 DEPRECATED
            schemaService.deprecate(previousActive.get());
        }

        // 守卫条件 3:依赖完整性
        dependencyChecker.validateDependencies(entity, worldId);

        // 状态转换
        entity.setStatus(SchemaStatus.ACTIVE);
        SchemaEntity activated = schemaService.update(entity);

        // 记录生命周期事件
        schemaService.recordLifecycleEvent(
            entity.getSchemaId(),
            SchemaLifecycleEvent.ACTIVATED,
            WorldContextHolder.getUserId());

        responseObserver.onNext(mapper.toResponse(activated));
        responseObserver.onCompleted();
    } catch (Exception e) {
        handleException(e, responseObserver);
    }
}

三重守卫条件

  1. 状态检查:只有 DRAFT 状态可以被激活
  2. 兼容性检查:如果同名 Schema 已有 ACTIVE 版本,新版本必须通过兼容性验证
  3. 依赖完整性:Schema 引用的其他 Schema(如 LinkType 引用的 ObjectType)必须存在且为 ACTIVE

自动弃用旧版本:激活新版本时,旧的 ACTIVE 版本自动转为 DEPRECATED。这保证了同一时刻只有一个 ACTIVE 版本,避免歧义。

#5. 兼容性检查:CompatibilityChecker 三级策略

Java
@Component
public class CompatibilityChecker {

    public CompatibilityResult check(
            SchemaEntity existing, SchemaEntity proposed) {
        CompatibilityMode mode = existing.getCompatibilityMode();
        List<String> violations = new ArrayList<>();

        switch (mode) {
            case BACKWARD:
                // 新版本可以读取旧版本写入的数据
                checkBackwardCompatibility(existing, proposed, violations);
                break;
            case FORWARD:
                // 旧版本可以读取新版本写入的数据
                checkForwardCompatibility(existing, proposed, violations);
                break;
            case FULL:
                // 双向兼容
                checkBackwardCompatibility(existing, proposed, violations);
                checkForwardCompatibility(existing, proposed, violations);
                break;
            case NONE:
                // 不检查
                break;
        }

        return new CompatibilityResult(violations.isEmpty(), violations);
    }

    private void checkBackwardCompatibility(
            SchemaEntity existing, SchemaEntity proposed,
            List<String> violations) {
        Map<String, PropertyDef> existingProps = parseProperties(existing);
        Map<String, PropertyDef> proposedProps = parseProperties(proposed);

        // 规则 1:不可删除已有的 required 字段
        for (Map.Entry<String, PropertyDef> entry :
                existingProps.entrySet()) {
            if (entry.getValue().isRequired()
                    && !proposedProps.containsKey(entry.getKey())) {
                violations.add(String.format(
                    "Cannot remove required property '%s' "
                    + "(backward incompatible)", entry.getKey()));
            }
        }

        // 规则 2:不可缩小字段类型范围
        for (Map.Entry<String, PropertyDef> entry :
                proposedProps.entrySet()) {
            PropertyDef existingProp =
                existingProps.get(entry.getKey());
            if (existingProp != null
                    && !isTypeWidening(
                        existingProp.getType(),
                        entry.getValue().getType())) {
                violations.add(String.format(
                    "Cannot narrow type of property '%s' from %s to %s",
                    entry.getKey(), existingProp.getType(),
                    entry.getValue().getType()));
            }
        }

        // 规则 3:新增的 required 字段必须有默认值
        for (Map.Entry<String, PropertyDef> entry :
                proposedProps.entrySet()) {
            if (!existingProps.containsKey(entry.getKey())
                    && entry.getValue().isRequired()
                    && !entry.getValue().hasDefaultValue()) {
                violations.add(String.format(
                    "New required property '%s' must have a default value",
                    entry.getKey()));
            }
        }
    }

    private void checkForwardCompatibility(
            SchemaEntity existing, SchemaEntity proposed,
            List<String> violations) {
        Map<String, PropertyDef> existingProps = parseProperties(existing);
        Map<String, PropertyDef> proposedProps = parseProperties(proposed);

        // 前向兼容规则:不可添加新的 required 字段(旧版本无法理解)
        for (Map.Entry<String, PropertyDef> entry :
                proposedProps.entrySet()) {
            if (!existingProps.containsKey(entry.getKey())
                    && entry.getValue().isRequired()) {
                violations.add(String.format(
                    "Cannot add required property '%s' "
                    + "(forward incompatible)", entry.getKey()));
            }
        }
    }

    private boolean isTypeWidening(String from, String to) {
        // int → long:允许(拓宽)
        // long → int:禁止(缩窄)
        // string → string:允许(同类型)
        Map<String, Integer> typeWidth = Map.of(
            "boolean", 1,
            "int", 2, "integer", 2,
            "long", 3,
            "float", 4,
            "double", 5,
            "string", 10  // string 可以容纳任何类型
        );
        return typeWidth.getOrDefault(to.toLowerCase(), 0)
            >= typeWidth.getOrDefault(from.toLowerCase(), 0);
    }
}

三级兼容策略

模式含义适用场景
BACKWARD新消费者能读旧生产者数据大多数场景
FORWARD旧消费者能读新生产者数据渐进式部署
FULL双向兼容严格环境
NONE不检查开发阶段

这种设计借鉴了 Apache Avro / Confluent Schema Registry 的兼容性模型,但适配了本体属性的语义。

#6. 依赖分析:SchemaDependencyChecker

Java
@Component
public class SchemaDependencyChecker {

    private final DependencyAnalyzer analyzer;
    private final DependencyCacheService cacheService;

    public void validateDependencies(
            SchemaEntity schema, String worldId) {
        DependencyGraph graph =
            cacheService.getOrBuild(worldId, () ->
                analyzer.buildGraph(worldId));

        // 检查所有引用的 Schema 是否存在且为 ACTIVE
        List<SchemaEdge> edges = graph.getOutgoingEdges(
            schema.getSchemaId());
        for (SchemaEdge edge : edges) {
            SchemaNode target = graph.getNode(edge.getTargetId());
            if (target == null) {
                throw new DependencyMissingException(
                    schema.getSchemaName(),
                    edge.getTargetId());
            }
            if (target.getStatus() != SchemaStatus.ACTIVE) {
                throw new DependencyNotActiveException(
                    schema.getSchemaName(),
                    target.getSchemaName(),
                    target.getStatus());
            }
        }
    }

    public ImpactAnalysis analyzeImpact(
            String schemaId, String worldId) {
        DependencyGraph graph =
            cacheService.getOrBuild(worldId, () ->
                analyzer.buildGraph(worldId));

        // BFS 计算所有受影响的下游 Schema
        List<ImpactItem> impacted = new ArrayList<>();
        Queue<String> queue = new LinkedList<>();
        Set<String> visited = new HashSet<>();
        queue.add(schemaId);

        while (!queue.isEmpty()) {
            String current = queue.poll();
            if (visited.contains(current)) continue;
            visited.add(current);

            List<SchemaEdge> incoming =
                graph.getIncomingEdges(current);
            for (SchemaEdge edge : incoming) {
                ImpactLevel level = (edge.getDependencyType()
                    == DependencyType.HARD)
                    ? ImpactLevel.BREAKING
                    : ImpactLevel.WARNING;
                impacted.add(new ImpactItem(
                    edge.getSourceId(),
                    graph.getNode(edge.getSourceId()).getSchemaName(),
                    level));
                queue.add(edge.getSourceId());
            }
        }

        return new ImpactAnalysis(schemaId, impacted);
    }
}

依赖图是延迟构建、按 World 缓存的。每个 World 有独立的 Schema 空间,因此依赖图也是隔离的。DependencyCacheService 在 Schema 发生变更时失效缓存。

**影响分析(Impact Analysis)**使用 BFS 遍历依赖图的逆向边,找出所有会被影响的下游 Schema。这个功能在弃用 Schema 之前非常有用——可以提前知道哪些下游会受影响。

#7. 版本管理:SchemaVersionHistory

Java
public class SchemaVersionHistory {
    private final List<VersionEntry> entries;

    @Value
    public static class VersionEntry {
        int version;
        String schemaId;
        SchemaStatus status;
        String changedBy;
        LocalDateTime changedAt;
        String changeDescription;
        String checksum;         // SHA-256 of definition
    }

    public void addVersion(SchemaEntity entity,
            String changeDescription) {
        String checksum = DigestUtils.sha256Hex(
            entity.getDefinition());

        // 防止重复提交相同内容
        if (!entries.isEmpty()) {
            VersionEntry latest = entries.get(entries.size() - 1);
            if (latest.getChecksum().equals(checksum)) {
                throw new DuplicateVersionException(
                    "Schema content unchanged, "
                    + "version bump not needed");
            }
        }

        entries.add(new VersionEntry(
            entity.getVersion(),
            entity.getSchemaId(),
            entity.getStatus(),
            entity.getUpdatedBy(),
            entity.getUpdatedAt(),
            changeDescription,
            checksum
        ));
    }
}

版本链是不可变的——只有 addVersion 操作,没有删除或修改操作。每个版本记录包含 SHA-256 校验和,用于检测重复提交。如果定义内容完全相同,拒绝创建新版本,避免版本号膨胀。

#8. Schema CRUD 完整流程

#8.1 更新 Schema

Java
@Override
public void updateSchema(UpdateSchemaRequest request,
        StreamObserver<SchemaResponse> responseObserver) {
    try {
        String worldId = WorldContextHolder.getWorldId();
        String userId = WorldContextHolder.getUserId();

        SchemaEntity existing = schemaService.getById(
            request.getSchemaId(), worldId);

        // ACTIVE 状态下修改需要额外验证
        if (existing.getStatus() == SchemaStatus.ACTIVE) {
            // 创建新的 DRAFT 版本而不是直接修改
            SchemaEntity newVersion = mapper.toDomain(
                request.getSchema(), worldId, userId);
            newVersion.setVersion(existing.getVersion() + 1);
            newVersion.setStatus(SchemaStatus.DRAFT);
            SchemaEntity saved = schemaService.register(newVersion);
            responseObserver.onNext(mapper.toResponse(saved));
        } else if (existing.getStatus() == SchemaStatus.DRAFT) {
            // DRAFT 可以直接修改
            mapper.updateEntity(existing, request.getSchema());
            existing.setUpdatedBy(userId);
            existing.setUpdatedAt(LocalDateTime.now());
            SchemaEntity updated = schemaService.update(existing);
            responseObserver.onNext(mapper.toResponse(updated));
        } else {
            throw Status.FAILED_PRECONDITION
                .withDescription("Cannot update schema in "
                    + existing.getStatus() + " state")
                .asRuntimeException();
        }

        responseObserver.onCompleted();
    } catch (Exception e) {
        handleException(e, responseObserver);
    }
}

关键设计:ACTIVE 状态的 Schema 不能直接修改——修改操作会创建一个新的 DRAFT 版本(版本号 +1)。这保证了生产环境的稳定性:正在运行的服务读取的是 ACTIVE 版本,新版本需要通过 activate() 流程(含兼容性检查)才能替代。

#8.2 弃用与归档

Java
@Override
public void deprecateSchema(DeprecateSchemaRequest request,
        StreamObserver<SchemaResponse> responseObserver) {
    try {
        String worldId = WorldContextHolder.getWorldId();
        SchemaEntity entity = schemaService.getById(
            request.getSchemaId(), worldId);

        if (entity.getStatus() != SchemaStatus.ACTIVE) {
            throw Status.FAILED_PRECONDITION
                .withDescription("Only ACTIVE schemas can be deprecated")
                .asRuntimeException();
        }

        // 影响分析
        ImpactAnalysis impact = dependencyChecker.analyzeImpact(
            entity.getSchemaId(), worldId);
        if (impact.hasBreakingImpact() && !request.getForce()) {
            throw Status.FAILED_PRECONDITION
                .withDescription("Schema has breaking dependencies: "
                    + impact.getSummary()
                    + ". Use force=true to override")
                .asRuntimeException();
        }

        entity.setStatus(SchemaStatus.DEPRECATED);
        SchemaEntity deprecated = schemaService.update(entity);

        responseObserver.onNext(mapper.toResponse(deprecated));
        responseObserver.onCompleted();
    } catch (Exception e) {
        handleException(e, responseObserver);
    }
}

弃用操作会先执行影响分析。如果有 BREAKING 级别的下游依赖,除非显式传入 force=true,否则拒绝弃用。

#9. 批量导入导出:YAML/JSON 双格式

Java
@Override
public void exportSchemas(ExportSchemasRequest request,
        StreamObserver<ExportSchemasResponse> responseObserver) {
    try {
        String worldId = WorldContextHolder.getWorldId();
        List<SchemaEntity> schemas = schemaService.listAll(worldId);

        String format = request.getFormat().isEmpty()
            ? "yaml" : request.getFormat();

        ObjectMapper mapper;
        if ("yaml".equalsIgnoreCase(format)) {
            mapper = new ObjectMapper(new YAMLFactory());
        } else {
            mapper = new ObjectMapper();
        }

        List<Map<String, Object>> exportData = schemas.stream()
            .map(this::toExportMap)
            .collect(Collectors.toList());

        byte[] content = mapper.writeValueAsBytes(exportData);

        responseObserver.onNext(
            ExportSchemasResponse.newBuilder()
                .setContent(ByteString.copyFrom(content))
                .setFormat(format)
                .setSchemaCount(schemas.size())
                .build());
        responseObserver.onCompleted();
    } catch (Exception e) {
        handleException(e, responseObserver);
    }
}

双格式设计:YAML 用于人类可读的版本控制(可以 commit 到 Git),JSON 用于程序间传输。导出包含完整的 Schema 定义、状态、版本历史和兼容性模式。

#10. EventRef 验证与跨实体引用

Java
@Component
public class EventRefValidator {

    private final SchemaService schemaService;

    public void validate(EventRefProto eventRef, String worldId) {
        // 验证引用的 ObjectType 存在
        String objectTypeId = eventRef.getObjectTypeId();
        Optional<SchemaEntity> objectType =
            schemaService.findActiveBySchemaId(objectTypeId, worldId);

        if (objectType.isEmpty()) {
            throw Status.FAILED_PRECONDITION
                .withDescription(String.format(
                    "Referenced ObjectType '%s' not found or not ACTIVE",
                    objectTypeId))
                .asRuntimeException();
        }

        // 验证引用的属性存在于 ObjectType 定义中
        for (String propertyRef : eventRef.getPropertyRefsList()) {
            if (!hasProperty(objectType.get(), propertyRef)) {
                throw Status.INVALID_ARGUMENT
                    .withDescription(String.format(
                        "Property '%s' not found in ObjectType '%s'",
                        propertyRef, objectTypeId))
                    .asRuntimeException();
            }
        }
    }
}

EventRef 是 coomia-dip 的 Object-Event 模型的核心概念:每个 ObjectType 可以关联事件源,事件的属性引用必须指向 ObjectType 中已定义的属性。EventRefValidator 确保这种引用的完整性。

#11. Key Takeaways

  1. 四态状态机:DRAFT → ACTIVE → DEPRECATED → ARCHIVED,每次转换都有明确的守卫条件
  2. ACTIVE 不可变原则:修改 ACTIVE Schema 会创建新 DRAFT 版本,保护生产环境稳定
  3. 三级兼容性模型:BACKWARD / FORWARD / FULL,借鉴 Schema Registry 最佳实践
  4. 依赖图 + 影响分析:弃用前自动评估下游影响,支持 force override
  5. 版本链不可变:SHA-256 校验防重复,只增不改
  6. YAML/JSON 双格式导出:支持 Git 版本控制和程序间传输
  7. EventRef 引用完整性:跨实体引用在注册时即验证,防止运行时 NPE

#下一篇

S9-03:WorldManagerService — 数据世界的 Git 操作。我们将深入 Nessie 集成层,看 coomia-dip 如何用 Git 语义管理数据分支、合并、发布和时间旅行。

Tags: #coomia-dip #source-code-reading #control-Layer #spring-boot #grpc #state-machine #schema-registry #compatibility-check #dependency-analysis