源码精读:OntologyRuntimeService — 实体 CRUD 的核心实现
OntologyRuntimeService 是 coomia-dip 数据层(Data Layer)中最核心的服务,承担着本体实例的全生命周期管理。它基于 Quarkus 3.x 框架,通过 gRPC 协议暴露接口,支持单条与批量的 CRUD 操作,内置软删除、版本历史、变更事件发射等企业级特性。本文将逐行剖析其架构分层、WorldContext 解析策略、Proto-Domain 双向转换、批量操作的错误聚合模式,以及变更触发器的调度机制。
源码精读:OntologyRuntimeService — 实体 CRUD 的核心实现
“系列:S9 源码精读 · 第 1 篇 | 难度:高级 | 阅读时间:25 分钟
#TL;DR
OntologyRuntimeService 是 coomia-dip 数据层(Data Layer)中最核心的服务,承担着本体实例的全生命周期管理。它基于 Quarkus 3.x 框架,通过 gRPC 协议暴露接口,支持单条与批量的 CRUD 操作,内置软删除、版本历史、变更事件发射等企业级特性。本文将逐行剖析其架构分层、WorldContext 解析策略、Proto-Domain 双向转换、批量操作的错误聚合模式,以及变更触发器的调度机制。
#目录
- 整体架构与分层设计
- gRPC 服务层:OntologyRuntimeGrpcService
- WorldContext 解析:双源优先级策略
- Proto 与 Domain 的双向转换:ProtoConverter
- 单实体 CRUD 实现
- 批量操作与部分失败语义
- 软删除与历史版本
- 关系操作:CreateRelation / DeleteRelation / GetRelations
- 变更事件触发器:OntologyChangeTriggerDispatcher
- 异常处理:GrpcExceptionHandler 统一映射
- Key Takeaways
#1. 整体架构与分层设计
OntologyRuntimeService 的代码分布在三个包中,遵循经典的六边形架构:
data-Layer/src/main/java/com/onto/data/
├── api/grpc/ # 入站适配器(gRPC)
│ ├── OntologyRuntimeGrpcService.java # gRPC 服务实现
│ ├── converter/ProtoConverter.java # Proto ↔ Domain 转换
│ ├── exception/GrpcExceptionHandler.java
│ └── interceptor/WorldContextInterceptor.java
├── service/ # 应用服务层
│ ├── OntologyRuntimeService.java # 接口定义
│ ├── DefaultOntologyRuntimeService.java
│ └── dto/ # 数据传输对象
│ ├── BatchCreateResponse.java
│ ├── BatchUpdateResponse.java
│ ├── UpdateMode.java
│ └── UpdateOptions.java
├── domain/instance/ # 领域模型
│ └── OntologyInstance.java
├── repository/ # 出站适配器(存储)
│ └── exception/InstanceNotFoundException.java
└── event/ # 事件机制
└── OntologyChangeTriggerDispatcher.java
设计决策:gRPC 服务层不包含业务逻辑——它只负责三件事:(1) 提取 WorldContext,(2) Proto ↔ Domain 转换,(3) 异常映射。业务规则全部下沉到 OntologyRuntimeService 接口及其实现类。
这种分层的好处在于,当我们需要增加 Flight SQL 或 REST 入口时,只需要新增适配器,核心服务层零改动。
#2. gRPC 服务层:OntologyRuntimeGrpcService
@GrpcService
@Blocking
public class OntologyRuntimeGrpcService
extends OntologyRuntimeServiceGrpc.OntologyRuntimeServiceImplBase {
private static final Logger LOG = LoggerFactory.getLogger(OntologyRuntimeGrpcService.class);
private final OntologyRuntimeService service;
private final ProtoConverter converter;
private final GrpcExceptionHandler exceptionHandler;
private final OntologyChangeTriggerDispatcher triggerDispatcher;
@Inject
public OntologyRuntimeGrpcService(
OntologyRuntimeService service,
ProtoConverter converter,
GrpcExceptionHandler exceptionHandler,
OntologyChangeTriggerDispatcher triggerDispatcher) {
this.service = service;
this.converter = converter;
this.exceptionHandler = exceptionHandler;
this.triggerDispatcher = triggerDispatcher;
}
}
关键注解解读:
@GrpcService:Quarkus gRPC 扩展的服务注册注解,自动将此类注册为 gRPC 服务端点@Blocking:声明所有方法在阻塞线程池中执行(而非 Vert.x 事件循环),因为底层涉及 JDBC/Doris 同步调用@Inject:CDI 构造器注入,四个依赖项各司其职
为什么使用 @Blocking? Quarkus 默认在 Vert.x IO 线程上处理 gRPC 请求,但 OntologyRuntimeService 的实现依赖 Doris JDBC 连接池(同步阻塞),如果在 IO 线程执行会导致 BlockingNotAllowedException。@Blocking 将请求调度到工作线程池,代价是上下文切换,但保证了线程安全。
#3. WorldContext 解析:双源优先级策略
WorldContext 是 coomia-dip 的核心隔离单元。每个请求必须携带 WorldContext,否则服务会拒绝处理。
private WorldContext resolveWorldContext(RequestContext requestContext) {
// 优先级 1:gRPC Interceptor 注入的 Context Key
WorldContext fromInterceptor = WorldContextInterceptor.WORLD_CONTEXT_KEY.get();
if (fromInterceptor != null) {
return fromInterceptor;
}
// 优先级 2:请求体中的 RequestContext.world 字段
if (requestContext != null && requestContext.hasWorld()) {
return WorldContext.of(
requestContext.getWorld().getWorldId(),
requestContext.getWorld().getBranch()
);
}
throw new IllegalStateException(
"WorldContext not found in gRPC context or request"
);
}
双源策略的设计考量:
| 来源 | 场景 | 优先级 |
|---|---|---|
| gRPC Interceptor | SDK 客户端自动注入 metadata | 高 |
| RequestContext 字段 | 手动构造请求、测试、跨网关调用 | 低 |
当两个来源都存在时,Interceptor 的值优先。这是因为 Interceptor 在请求到达服务方法之前执行,已经过认证和校验,安全性更高。
#4. Proto 与 Domain 的双向转换:ProtoConverter
ProtoConverter 是一个无状态的转换器,负责 Protobuf 消息与领域对象之间的双向映射:
@ApplicationScoped
public class ProtoConverter {
public OntologyInstance toDomain(CreateInstanceInput input, WorldContext ctx) {
OntologyInstance instance = new OntologyInstance();
instance.setEntityType(input.getEntityType());
instance.setPrimaryKey(input.hasPrimaryKey()
? input.getPrimaryKey()
: UUID.randomUUID().toString());
instance.setWorldId(ctx.getWorldId());
instance.setBranch(ctx.getBranch());
// 属性映射:Proto Struct → Java Map
Map<String, Object> attributes = new HashMap<>();
for (Map.Entry<String, Value> entry :
input.getAttributesMap().entrySet()) {
attributes.put(entry.getKey(),
AttributeValueConverter.fromProto(entry.getValue()));
}
instance.setAttributes(attributes);
return instance;
}
public InstanceProto toProto(OntologyInstance instance) {
InstanceProto.Builder builder = InstanceProto.newBuilder()
.setInstanceId(instance.getId())
.setEntityType(instance.getEntityType())
.setPrimaryKey(instance.getPrimaryKey())
.setVersion(instance.getVersion())
.setCreatedAt(toTimestamp(instance.getCreatedAt()))
.setUpdatedAt(toTimestamp(instance.getUpdatedAt()));
if (instance.isDeleted()) {
builder.setDeletedAt(toTimestamp(instance.getDeletedAt()));
}
// 属性映射:Java Map → Proto Struct
for (Map.Entry<String, Object> entry :
instance.getAttributes().entrySet()) {
builder.putAttributes(entry.getKey(),
AttributeValueConverter.toProto(entry.getValue()));
}
return builder.build();
}
}
设计亮点:
- 主键策略:如果客户端未提供
primaryKey,自动生成 UUID——这支持两种使用模式:业务键(如员工工号)和系统键 - AttributeValueConverter:处理 Proto
Value与 JavaObject之间的递归类型映射,支持 null、string、number、bool、list、map 六种类型 - 版本号传递:
toProto始终输出 version 字段,支持乐观锁场景
#5. 单实体 CRUD 实现
#5.1 创建实例
@Override
public void createInstance(CreateInstanceRequest request,
StreamObserver<CreateInstanceResponse> responseObserver) {
try {
WorldContext ctx = resolveWorldContext(request.getContext());
OntologyInstance domain = converter.toDomain(
request.getInstance(), ctx);
OntologyInstance created = service.createInstance(domain, ctx);
// 发射变更事件
triggerDispatcher.dispatch(
ChangeType.CREATE, created.getEntityType(),
created.getId(), ctx);
CreateInstanceResponse response = CreateInstanceResponse.newBuilder()
.setInstance(converter.toProto(created))
.build();
responseObserver.onNext(response);
responseObserver.onCompleted();
} catch (Exception e) {
exceptionHandler.handle(e, responseObserver);
}
}
执行流程:
客户端 → gRPC → resolveWorldContext → converter.toDomain
→ service.createInstance(业务逻辑 + 存储)
→ triggerDispatcher.dispatch(异步事件)
→ converter.toProto → responseObserver.onNext
注意 triggerDispatcher.dispatch 在成功创建之后、响应之前调用。这意味着事件发射失败不会导致创建回滚,但会记录错误日志。这是一个典型的"最终一致性"选择。
#5.2 获取实例
@Override
public void getInstance(GetInstanceRequest request,
StreamObserver<GetInstanceResponse> responseObserver) {
try {
WorldContext ctx = resolveWorldContext(request.getContext());
Optional<OntologyInstance> found = service.getInstance(
request.getEntityType(),
request.getInstanceId(),
ctx
);
if (found.isEmpty()) {
throw new InstanceNotFoundException(
request.getEntityType(),
request.getInstanceId()
);
}
GetInstanceResponse response = GetInstanceResponse.newBuilder()
.setInstance(converter.toProto(found.get()))
.build();
responseObserver.onNext(response);
responseObserver.onCompleted();
} catch (Exception e) {
exceptionHandler.handle(e, responseObserver);
}
}
注意:这里没有使用 Optional.orElseThrow(),而是显式检查 isEmpty()。这是因为 InstanceNotFoundException 需要传递 entityType 和 instanceId 两个参数用于错误消息构造,lambda 写法反而更冗长。
#5.3 更新实例
@Override
public void updateInstance(UpdateInstanceRequest request,
StreamObserver<UpdateInstanceResponse> responseObserver) {
try {
WorldContext ctx = resolveWorldContext(request.getContext());
UpdateInstanceInput input = request.getInstance();
UpdateOptions options = UpdateOptions.builder()
.mode(input.hasUpdateMode()
? UpdateMode.valueOf(input.getUpdateMode().name())
: UpdateMode.MERGE)
.expectedVersion(input.hasExpectedVersion()
? input.getExpectedVersion()
: null)
.build();
OntologyInstance updated = service.updateInstance(
input.getEntityType(),
input.getInstanceId(),
converter.toAttributeMap(input.getAttributesMap()),
options,
ctx
);
triggerDispatcher.dispatch(
ChangeType.UPDATE, updated.getEntityType(),
updated.getId(), ctx);
UpdateInstanceResponse response = UpdateInstanceResponse.newBuilder()
.setInstance(converter.toProto(updated))
.build();
responseObserver.onNext(response);
responseObserver.onCompleted();
} catch (Exception e) {
exceptionHandler.handle(e, responseObserver);
}
}
两种更新模式:
| 模式 | 行为 | 用途 |
|---|---|---|
MERGE | 仅更新传入的属性,保留其他属性 | 增量更新 |
REPLACE | 用传入属性完全替换现有属性 | 全量覆盖 |
乐观锁:expectedVersion 是可选的。当提供时,服务层会比较当前版本号,不匹配则抛出 OptimisticLockException,由 GrpcExceptionHandler 映射为 Status.ABORTED。
#6. 批量操作与部分失败语义
批量操作是 coomia-dip 的重要特性。与单条操作的"全部或失败"不同,批量操作采用部分成功语义:
@Override
public void batchCreateInstances(BatchCreateInstancesRequest request,
StreamObserver<BatchCreateInstancesResponse> responseObserver) {
try {
WorldContext ctx = resolveWorldContext(request.getContext());
List<OntologyInstance> domainList = request.getInstancesList()
.stream()
.map(input -> converter.toDomain(input, ctx))
.collect(Collectors.toList());
BatchCreateResponse batchResult =
service.batchCreateInstances(domainList, ctx);
BatchCreateInstancesResponse.Builder responseBuilder =
BatchCreateInstancesResponse.newBuilder()
.setTotalRequested(request.getInstancesCount())
.setSuccessCount(batchResult.getSuccessful().size())
.setFailureCount(batchResult.getErrors().size());
// 成功的实例
for (OntologyInstance created : batchResult.getSuccessful()) {
responseBuilder.addInstances(converter.toProto(created));
triggerDispatcher.dispatch(
ChangeType.CREATE, created.getEntityType(),
created.getId(), ctx);
}
// 失败的记录
for (BatchOperationError error : batchResult.getErrors()) {
responseBuilder.addErrors(
BatchError.newBuilder()
.setIndex(error.getIndex())
.setMessage(error.getMessage())
.setCode(error.getCode())
.build()
);
}
responseObserver.onNext(responseBuilder.build());
responseObserver.onCompleted();
} catch (Exception e) {
exceptionHandler.handle(e, responseObserver);
}
}
部分失败响应结构:
message BatchCreateInstancesResponse {
int32 total_requested = 1;
int32 success_count = 2;
int32 failure_count = 3;
repeated InstanceProto instances = 4; // 成功的
repeated BatchError errors = 5; // 失败的,携带 index
}
设计决策:每个 BatchError 包含原始请求中的 index,客户端可以据此定位哪些记录失败。这比"全部回滚"更适合大批量导入场景——成功 9999 条、失败 1 条,不需要全部重来。
#7. 软删除与历史版本
#7.1 软删除
@Override
public void deleteInstance(DeleteInstanceRequest request,
StreamObserver<DeleteInstanceResponse> responseObserver) {
try {
WorldContext ctx = resolveWorldContext(request.getContext());
boolean deleted = service.deleteInstance(
request.getEntityType(),
request.getInstanceId(),
request.getHardDelete(), // 是否物理删除
ctx
);
triggerDispatcher.dispatch(
ChangeType.DELETE,
request.getEntityType(),
request.getInstanceId(),
ctx);
DeleteInstanceResponse response = DeleteInstanceResponse.newBuilder()
.setDeleted(deleted)
.build();
responseObserver.onNext(response);
responseObserver.onCompleted();
} catch (Exception e) {
exceptionHandler.handle(e, responseObserver);
}
}
默认情况下,deleteInstance 执行软删除:设置 deleted_at 时间戳,实例在正常查询中不可见,但在审计查询和历史回溯中仍然可访问。只有显式传入 hardDelete=true 才会物理删除记录。
#7.2 历史版本查询
@Override
public void getInstanceHistory(GetInstanceHistoryRequest request,
StreamObserver<GetInstanceHistoryResponse> responseObserver) {
try {
WorldContext ctx = resolveWorldContext(request.getContext());
List<OntologyInstance> history = service.getInstanceHistory(
request.getEntityType(),
request.getInstanceId(),
request.getMaxVersions(),
ctx
);
GetInstanceHistoryResponse.Builder builder =
GetInstanceHistoryResponse.newBuilder();
for (OntologyInstance version : history) {
builder.addVersions(converter.toProto(version));
}
responseObserver.onNext(builder.build());
responseObserver.onCompleted();
} catch (Exception e) {
exceptionHandler.handle(e, responseObserver);
}
}
历史版本查询利用 Nessie 的 Git-like 版本管理能力,每次更新都会创建新的 commit。maxVersions 参数限制返回的版本数量,避免大对象的全量历史导致内存溢出。
还有一个按版本号精确获取的方法:
@Override
public void getInstanceAtVersion(GetInstanceAtVersionRequest request,
StreamObserver<GetInstanceResponse> responseObserver) {
try {
WorldContext ctx = resolveWorldContext(request.getContext());
Optional<OntologyInstance> found = service.getInstanceAtVersion(
request.getEntityType(),
request.getInstanceId(),
request.getVersion(),
ctx
);
// ...
} catch (Exception e) {
exceptionHandler.handle(e, responseObserver);
}
}
这支持"时间旅行"查询——查看实体在任意历史时刻的状态。
#8. 关系操作:CreateRelation / DeleteRelation / GetRelations
关系(Relation)是本体模型的核心概念之一。实体之间通过关系连接,形成知识图谱。
@Override
public void createRelation(CreateRelationRequest request,
StreamObserver<CreateRelationResponse> responseObserver) {
try {
WorldContext ctx = resolveWorldContext(request.getContext());
Relation relation = converter.toRelationDomain(
request.getRelation(), ctx);
Relation created = service.createRelation(relation, ctx);
triggerDispatcher.dispatch(
ChangeType.RELATION_CREATE,
relation.getRelationType(),
created.getId(), ctx);
responseObserver.onNext(
CreateRelationResponse.newBuilder()
.setRelation(converter.toRelationProto(created))
.build()
);
responseObserver.onCompleted();
} catch (Exception e) {
exceptionHandler.handle(e, responseObserver);
}
}
@Override
public void getRelations(GetRelationsRequest request,
StreamObserver<GetRelationsResponse> responseObserver) {
try {
WorldContext ctx = resolveWorldContext(request.getContext());
List<Relation> relations = service.getRelations(
request.getEntityType(),
request.getInstanceId(),
request.hasRelationType()
? request.getRelationType() : null,
request.hasDirection()
? Direction.valueOf(request.getDirection().name())
: Direction.BOTH,
ctx
);
GetRelationsResponse.Builder builder =
GetRelationsResponse.newBuilder();
for (Relation rel : relations) {
builder.addRelations(converter.toRelationProto(rel));
}
responseObserver.onNext(builder.build());
responseObserver.onCompleted();
} catch (Exception e) {
exceptionHandler.handle(e, responseObserver);
}
}
关系查询的方向性:Direction 枚举支持三种值:OUTGOING(从此实体出发)、INCOMING(指向此实体)、BOTH(双向)。底层存储在图数据库(TuGraph)中是有向边,但查询时可以反向遍历。
#9. 变更事件触发器:OntologyChangeTriggerDispatcher
@ApplicationScoped
public class OntologyChangeTriggerDispatcher {
@Inject
Event<OntologyChangeEvent> changeEventBus;
@Inject
SubscriptionEventRouter subscriptionRouter;
public void dispatch(ChangeType type, String entityType,
String instanceId, WorldContext ctx) {
OntologyChangeEvent event = OntologyChangeEvent.builder()
.changeType(type)
.entityType(entityType)
.instanceId(instanceId)
.worldId(ctx.getWorldId())
.branch(ctx.getBranch())
.timestamp(Instant.now())
.build();
// 1. CDI Event(进程内同步)
changeEventBus.fire(event);
// 2. Subscription Router(Kafka 异步)
subscriptionRouter.route(event);
LOG.debugf("Dispatched %s event for %s/%s",
type, entityType, instanceId);
}
}
事件分两路:
- CDI Event:进程内同步广播,用于触发派生属性重算、缓存失效等本地操作
- Subscription Router:将事件发送到 Kafka,由
SubscriptionService消费并投递给外部订阅者
这种双通道设计保证了本地副作用的即时性和跨进程通知的可靠性。
#10. 异常处理:GrpcExceptionHandler 统一映射
@ApplicationScoped
public class GrpcExceptionHandler {
private static final Map<Class<? extends Exception>, Status.Code> EXCEPTION_MAP =
Map.of(
InstanceNotFoundException.class, Status.Code.NOT_FOUND,
DuplicateKeyException.class, Status.Code.ALREADY_EXISTS,
OptimisticLockException.class, Status.Code.ABORTED,
SchemaValidationException.class, Status.Code.INVALID_ARGUMENT,
WorldContextMissingException.class, Status.Code.UNAUTHENTICATED,
QuotaExceededException.class, Status.Code.RESOURCE_EXHAUSTED
);
public <T> void handle(Exception e,
StreamObserver<T> responseObserver) {
Status.Code code = EXCEPTION_MAP.getOrDefault(
e.getClass(), Status.Code.INTERNAL);
LOG.errorf(e, "gRPC error [%s]: %s", code, e.getMessage());
responseObserver.onError(
Status.fromCode(code)
.withDescription(e.getMessage())
.withCause(e)
.asRuntimeException()
);
}
}
映射策略:使用静态 Map 而非 if-else 链,新增异常类型只需添加一行。未匹配的异常统一归为 INTERNAL,避免敏感信息泄露。
#11. Key Takeaways
- 六边形架构分层:gRPC 层只做协议适配(WorldContext 解析 + Proto 转换 + 异常映射),业务逻辑下沉到 Service 层
- WorldContext 双源解析:Interceptor 优先于 Request 字段,安全性与灵活性兼顾
- 批量操作部分成功:返回成功列表 + 错误列表(含 index),比全量回滚更适合数据导入场景
- 软删除默认:默认软删除保留审计能力,
hardDelete参数控制物理删除 - 变更事件双通道:CDI Event 保证本地即时性,Kafka 保证跨进程可靠投递
- 异常映射表:静态 Map 的声明式异常映射,可扩展性和可读性兼得
@Blocking注解:因为底层 Doris JDBC 是同步的,必须将 gRPC 处理调度到阻塞线程池
#下一篇
S9-02:SchemaRegistryService — 本体注册的状态机。我们将深入 Control Layer 的 Schema 注册中心,看它如何管理 Schema 的 DRAFT → ACTIVE → DEPRECATED 生命周期状态机,以及向前/向后兼容性检查的实现。
Tags: #coomia-dip #source-code-reading #data-Layer #quarkus #grpc #entity-crud #soft-delete #batch-operations #world-context