返回博客

源码精读:OntologyRuntimeService — 实体 CRUD 的核心实现

OntologyRuntimeService 是 coomia-dip 数据层(Data Layer)中最核心的服务,承担着本体实例的全生命周期管理。它基于 Quarkus 3.x 框架,通过 gRPC 协议暴露接口,支持单条与批量的 CRUD 操作,内置软删除、版本历史、变更事件发射等企业级特性。本文将逐行剖析其架构分层、WorldContext 解析策略、Proto-Domain 双向转换、批量操作的错误聚合模式,以及变更触发器的调度机制。

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

源码精读:OntologyRuntimeService — 实体 CRUD 的核心实现

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

#TL;DR

OntologyRuntimeService 是 coomia-dip 数据层(Data Layer)中最核心的服务,承担着本体实例的全生命周期管理。它基于 Quarkus 3.x 框架,通过 gRPC 协议暴露接口,支持单条与批量的 CRUD 操作,内置软删除、版本历史、变更事件发射等企业级特性。本文将逐行剖析其架构分层、WorldContext 解析策略、Proto-Domain 双向转换、批量操作的错误聚合模式,以及变更触发器的调度机制。

#目录

  1. 整体架构与分层设计
  2. gRPC 服务层:OntologyRuntimeGrpcService
  3. WorldContext 解析:双源优先级策略
  4. Proto 与 Domain 的双向转换:ProtoConverter
  5. 单实体 CRUD 实现
  6. 批量操作与部分失败语义
  7. 软删除与历史版本
  8. 关系操作:CreateRelation / DeleteRelation / GetRelations
  9. 变更事件触发器:OntologyChangeTriggerDispatcher
  10. 异常处理:GrpcExceptionHandler 统一映射
  11. Key Takeaways

#1. 整体架构与分层设计

OntologyRuntimeService 的代码分布在三个包中,遵循经典的六边形架构:

Code
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

Java
@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,否则服务会拒绝处理。

Java
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 InterceptorSDK 客户端自动注入 metadata
RequestContext 字段手动构造请求、测试、跨网关调用

当两个来源都存在时,Interceptor 的值优先。这是因为 Interceptor 在请求到达服务方法之前执行,已经过认证和校验,安全性更高。

#4. Proto 与 Domain 的双向转换:ProtoConverter

ProtoConverter 是一个无状态的转换器,负责 Protobuf 消息与领域对象之间的双向映射:

Java
@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();
    }
}

设计亮点

  1. 主键策略:如果客户端未提供 primaryKey,自动生成 UUID——这支持两种使用模式:业务键(如员工工号)和系统键
  2. AttributeValueConverter:处理 Proto Value 与 Java Object 之间的递归类型映射,支持 null、string、number、bool、list、map 六种类型
  3. 版本号传递toProto 始终输出 version 字段,支持乐观锁场景

#5. 单实体 CRUD 实现

#5.1 创建实例

Java
@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);
    }
}

执行流程

Code
客户端 → gRPC → resolveWorldContext → converter.toDomain
  → service.createInstance(业务逻辑 + 存储)
  → triggerDispatcher.dispatch(异步事件)
  → converter.toProto → responseObserver.onNext

注意 triggerDispatcher.dispatch 在成功创建之后、响应之前调用。这意味着事件发射失败不会导致创建回滚,但会记录错误日志。这是一个典型的"最终一致性"选择。

#5.2 获取实例

Java
@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 需要传递 entityTypeinstanceId 两个参数用于错误消息构造,lambda 写法反而更冗长。

#5.3 更新实例

Java
@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 的重要特性。与单条操作的"全部或失败"不同,批量操作采用部分成功语义:

Java
@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);
    }
}

部分失败响应结构

PROTOBUF
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 软删除

Java
@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 历史版本查询

Java
@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 参数限制返回的版本数量,避免大对象的全量历史导致内存溢出。

还有一个按版本号精确获取的方法:

Java
@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)是本体模型的核心概念之一。实体之间通过关系连接,形成知识图谱。

Java
@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

Java
@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);
    }
}

事件分两路:

  1. CDI Event:进程内同步广播,用于触发派生属性重算、缓存失效等本地操作
  2. Subscription Router:将事件发送到 Kafka,由 SubscriptionService 消费并投递给外部订阅者

这种双通道设计保证了本地副作用的即时性和跨进程通知的可靠性。

#10. 异常处理:GrpcExceptionHandler 统一映射

Java
@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

  1. 六边形架构分层:gRPC 层只做协议适配(WorldContext 解析 + Proto 转换 + 异常映射),业务逻辑下沉到 Service 层
  2. WorldContext 双源解析:Interceptor 优先于 Request 字段,安全性与灵活性兼顾
  3. 批量操作部分成功:返回成功列表 + 错误列表(含 index),比全量回滚更适合数据导入场景
  4. 软删除默认:默认软删除保留审计能力,hardDelete 参数控制物理删除
  5. 变更事件双通道:CDI Event 保证本地即时性,Kafka 保证跨进程可靠投递
  6. 异常映射表:静态 Map 的声明式异常映射,可扩展性和可读性兼得
  7. @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