源码精读:SchemaRegistryService — 本体注册的状态机
SchemaRegistryService 是 Control Layer(Control Layer)中管理本体 Schema 全生命周期的核心服务。它基于 Spring Boot 3.x + gRPC 构建,实现了 DRAFT → ACTIVE → DEPRECATED → ARCHIVED 的四态状态机,内置向前/向后兼容性检查、依赖分析、版本历史管理和批量导入导出。本文将深入分析状态机的转换规则、CompatibilityChecker 的三级兼容策略、SchemaDependencyChecker 的级联影响分析,以及 SchemaVersionHistory 的不可变版本链设计。
源码精读:SchemaRegistryService — 本体注册的状态机
“系列:S9 源码精读 · 第 2 篇 | 难度:高级 | 阅读时间:25 分钟
#TL;DR
SchemaRegistryService 是 Control Layer(Control Layer)中管理本体 Schema 全生命周期的核心服务。它基于 Spring Boot 3.x + gRPC 构建,实现了 DRAFT → ACTIVE → DEPRECATED → ARCHIVED 的四态状态机,内置向前/向后兼容性检查、依赖分析、版本历史管理和批量导入导出。本文将深入分析状态机的转换规则、CompatibilityChecker 的三级兼容策略、SchemaDependencyChecker 的级联影响分析,以及 SchemaVersionHistory 的不可变版本链设计。
#目录
- 架构定位与模块结构
- 领域模型:SchemaEntity 与 SchemaStatus
- gRPC 服务入口:SchemaRegistryServiceImpl
- 状态机:四态转换与守卫条件
- 兼容性检查:CompatibilityChecker 三级策略
- 依赖分析:SchemaDependencyChecker
- 版本管理:SchemaVersionHistory
- Schema CRUD 完整流程
- 批量导入导出:YAML/JSON 双格式
- EventRef 验证与跨实体引用
- Key Takeaways
#1. 架构定位与模块结构
Schema Registry 位于 Control Layer,是整个平台的"元数据大脑"。所有 Layer 在操作数据之前,都需要先向 Schema Registry 查询实体类型定义。
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
@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;
}
public enum SchemaStatus {
DRAFT, // 草稿,可自由修改
ACTIVE, // 已激活,修改需要兼容性检查
DEPRECATED, // 已弃用,只读
ARCHIVED // 已归档,不可见
}
状态转换图:
DRAFT ──activate()──→ ACTIVE ──deprecate()──→ DEPRECATED ──archive()──→ ARCHIVED
↑ │ │
│ │ revert() │ revert()
└─────────────────────┘ │
↑ │
└──────────────────────────────────────────────┘
设计考量:DRAFT 状态允许自由修改而不触发兼容性检查。只有从 DRAFT → ACTIVE 的 activate() 转换会执行完整的验证链。这使得 Schema 的迭代开发更高效——在 DRAFT 阶段可以反复调整字段。
#3. gRPC 服务入口:SchemaRegistryServiceImpl
@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. 状态机:四态转换与守卫条件
@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);
}
}
三重守卫条件:
- 状态检查:只有 DRAFT 状态可以被激活
- 兼容性检查:如果同名 Schema 已有 ACTIVE 版本,新版本必须通过兼容性验证
- 依赖完整性:Schema 引用的其他 Schema(如 LinkType 引用的 ObjectType)必须存在且为 ACTIVE
自动弃用旧版本:激活新版本时,旧的 ACTIVE 版本自动转为 DEPRECATED。这保证了同一时刻只有一个 ACTIVE 版本,避免歧义。
#5. 兼容性检查:CompatibilityChecker 三级策略
@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
@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
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
@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 弃用与归档
@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 双格式
@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 验证与跨实体引用
@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
- 四态状态机:DRAFT → ACTIVE → DEPRECATED → ARCHIVED,每次转换都有明确的守卫条件
- ACTIVE 不可变原则:修改 ACTIVE Schema 会创建新 DRAFT 版本,保护生产环境稳定
- 三级兼容性模型:BACKWARD / FORWARD / FULL,借鉴 Schema Registry 最佳实践
- 依赖图 + 影响分析:弃用前自动评估下游影响,支持 force override
- 版本链不可变:SHA-256 校验防重复,只增不改
- YAML/JSON 双格式导出:支持 Git 版本控制和程序间传输
- 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