源码精读:PipelineService — DSL 到 DAG 的编译
PipelineService 是 coomia-dip 数据管道层(Pipeline & Orchestration Layer,合并至 Data Layer)的核心服务,负责将声明式的 Pipeline DSL 编译为可执行的有向无环图(DAG)并调度执行。基于 Quarkus 3.x 框架,通过 gRPC 暴露 CRUD、执行、监控三组 API。本文将深入剖析 PipelineServiceGrpcImpl 适配器模式、PipelineExecutionService 应用服务层的 Pipeline 定义管理、Quarkus CDI 代理与 Mutiny Uni 的作用域陷阱、以及 Proto 合约中 PipelineDefinition 的 Step-DAG 编译模型。
源码精读:PipelineService — DSL 到 DAG 的编译
“系列:S9 源码精读 · 第 16 篇 | 难度:高级 | 阅读时间:25 分钟
#TL;DR
PipelineService 是 coomia-dip 数据管道层(Pipeline & Orchestration Layer,合并至 Data Layer)的核心服务,负责将声明式的 Pipeline DSL 编译为可执行的有向无环图(DAG)并调度执行。基于 Quarkus 3.x 框架,通过 gRPC 暴露 CRUD、执行、监控三组 API。本文将深入剖析 PipelineServiceGrpcImpl 适配器模式、PipelineExecutionService 应用服务层的 Pipeline 定义管理、Quarkus CDI 代理与 Mutiny Uni 的作用域陷阱、以及 Proto 合约中 PipelineDefinition 的 Step-DAG 编译模型。
#目录
- 整体架构:Pipeline & Orchestration Layer 的 Quarkus gRPC 服务
- Proto 合约:PipelineDefinition 与 Step-DAG 模型
- PipelineServiceGrpcImpl:适配器模式
- Quarkus CDI 代理的作用域陷阱
- CreatePipeline:定义与验证
- Pipeline 执行:DAG 调度策略
- 血缘记录:Pipeline 与 LineageService 的集成
- 错误处理:gRPC Status 映射
- Metadata 与版本管理
- Pipeline 监控与可观测性
- Key Takeaways
#1. 整体架构:Pipeline & Orchestration Layer 的 Quarkus gRPC 服务
Pipeline 服务位于 Data Layer(Data Layer + Pipeline & Orchestration Layer 合并后),使用 Quarkus 3.x 框架构建:
data-Layer/src/main/java/com/onto/plane_f/
├── api/grpc/
│ └── PipelineServiceGrpcImpl.java # gRPC 适配层
├── service/
│ └── PipelineExecutionService.java # 应用服务层
proto/plane_f/
├── pipeline_engine.proto # Pipeline 核心合约
└── lineage.proto # 血缘追踪合约
技术选型决策:Pipeline 服务选择 Quarkus 而非 Spring Boot,是因为 Data Layer 需要极低的启动时间(GraalVM native image 支持)和更小的内存占用。同时 Quarkus 的 Mutiny 响应式编程模型与 Pipeline 的异步执行天然契合。
#2. Proto 合约:PipelineDefinition 与 Step-DAG 模型
Pipeline 的 Proto 定义遵循 Step-DAG 模型——每个 Pipeline 由一组有依赖关系的 Step 组成:
message PipelineDefinition {
string pipeline_id = 1;
string name = 2;
string version = 3;
repeated PipelineStep steps = 4;
PipelineMetadata metadata = 5;
PipelineConfig config = 6;
}
message PipelineStep {
string step_id = 1;
string step_name = 2;
string step_type = 3; // EXTRACT, TRANSFORM, LOAD, CUSTOM
repeated string depends_on = 4; // 上游依赖的 step_id 列表
google.protobuf.Struct config = 5;
}
depends_on 字段将 Step 列表编译为 DAG——当一个 Step 的所有依赖都完成后,它才被调度执行。编译器会检测循环依赖并拒绝包含环的 Pipeline 定义。
Step 类型的设计遵循 ETL 范式:
- EXTRACT:从外部数据源提取数据
- TRANSFORM:数据转换和清洗
- LOAD:将处理后的数据加载到目标存储
- CUSTOM:用户自定义步骤,可引用 FunctionRuntime 中的函数
#3. PipelineServiceGrpcImpl:适配器模式
PipelineServiceGrpcImpl 是一个典型的薄适配器,不包含业务逻辑:
@GrpcService
public class PipelineServiceGrpcImpl implements PipelineService {
@Inject
PipelineExecutionService service;
@Override
public Uni<PipelineResponse> createPipeline(CreatePipelineRequest request) {
LOG.infof("CreatePipeline: request received");
try {
PipelineDefinition inputDef = request.getPipeline();
if (inputDef == null) {
return Uni.createFrom().failure(Status.INVALID_ARGUMENT
.withDescription("pipeline definition is required")
.asRuntimeException());
}
PipelineDefinition storedDef = service.createPipeline(inputDef);
Timestamp createdAt = storedDef.hasMetadata()
? storedDef.getMetadata().getCreatedAt()
: Timestamp.getDefaultInstance();
PipelineResponse response = PipelineResponse.newBuilder()
.setPipeline(storedDef)
.setCreatedAt(createdAt)
.setUpdatedAt(updatedAt)
.setVersionId(versionId)
.build();
return Uni.createFrom().item(response);
} catch (PipelineAlreadyExistsException e) {
return Uni.createFrom().failure(Status.ALREADY_EXISTS
.withDescription(e.getMessage()).asRuntimeException());
}
}
}
关键模式:所有业务逻辑调用都在方法体中同步执行,然后将结果包装在 Uni.createFrom().item(value) 中返回。这是一个刻意的设计决策——下一节将解释原因。
#4. Quarkus CDI 代理的作用域陷阱
代码中有一段重要的注释揭示了 Quarkus CDI 的陷阱:
/**
* IMPORTANT: All service calls are executed synchronously in the method body
* (NOT inside Uni lambdas) to avoid Quarkus CDI proxy scoping issues where
* Uni.createFrom().item(Supplier) lambdas resolve to different bean instances,
* causing the in-memory pipeline store to appear empty.
* The result is then wrapped with Uni.createFrom().item(value).
*/
当使用 Uni.createFrom().item(() -> service.createPipeline(inputDef)) 时,lambda 内部的 service 引用可能解析到一个不同的 CDI 代理实例——特别是当 PipelineExecutionService 使用 @ApplicationScoped 作用域时,lambda 可能在不同的线程上执行,导致看到的是一个"空的"内存存储。
解决方案:先同步调用 service.createPipeline(inputDef) 获取结果,再将结果(而非 Supplier)传给 Uni.createFrom().item(value)。这是一个经典的 Quarkus + CDI + Mutiny 集成陷阱。
#5. CreatePipeline:定义与验证
Pipeline 创建流程包含多层验证:
@Override
public Uni<PipelineResponse> createPipeline(CreatePipelineRequest request) {
try {
PipelineDefinition inputDef = request.getPipeline();
if (inputDef == null) {
return Uni.createFrom().failure(Status.INVALID_ARGUMENT
.withDescription("pipeline definition is required")
.asRuntimeException());
}
PipelineDefinition storedDef = service.createPipeline(inputDef);
// ... 构建响应 ...
} catch (PipelineAlreadyExistsException e) {
return Uni.createFrom().failure(Status.ALREADY_EXISTS
.withDescription(e.getMessage()).asRuntimeException());
} catch (IllegalArgumentException e) {
return Uni.createFrom().failure(Status.INVALID_ARGUMENT
.withDescription(e.getMessage()).asRuntimeException());
} catch (Exception e) {
LOG.errorf(e, "CreatePipeline failed");
return Uni.createFrom().failure(Status.INTERNAL
.withDescription("Failed to create pipeline: " + e.getMessage())
.asRuntimeException());
}
}
三层异常处理形成了清晰的错误分类:
PipelineAlreadyExistsException→ALREADY_EXISTS:Pipeline 名称唯一性冲突IllegalArgumentException→INVALID_ARGUMENT:输入验证失败(如 DAG 包含环)Exception→INTERNAL:未预期的内部错误
#6. Pipeline 执行:DAG 调度策略
Pipeline 执行使用拓扑排序确定 Step 的执行顺序:
- 拓扑排序:解析
depends_on构建 DAG,进行拓扑排序 - 并行调度:同一拓扑层级(无依赖关系)的 Step 可并行执行
- 失败处理:某个 Step 失败时,其所有下游 Step 自动标记为 SKIPPED
- 重试机制:支持 Step 级别的重试策略
DAG 的编译发生在 Pipeline 创建时——PipelineExecutionService.createPipeline() 会验证 DAG 的合法性(无环、所有引用的 step_id 存在)。运行时只需按照拓扑序调度。
#7. 血缘记录:Pipeline 与 LineageService 的集成
Pipeline 执行过程中会自动记录血缘信息到 Pipeline & Orchestration Layer 的 LineageService:
message RecordLineageRequest {
ExecutionContext context = 1;
repeated DatasetRef inputs = 2;
repeated DatasetRef outputs = 3;
repeated FieldMapping field_mappings = 4;
map<string, string> metadata = 5;
}
每个 Step 完成后,记录其输入和输出的数据集引用,以及字段级别的映射关系。这为事后的数据追溯和影响分析提供了完整的链路信息。
FieldMapping 支持六种转换类型:
enum TransformType {
TRANSFORM_TYPE_DIRECT = 1; // 直接复制
TRANSFORM_TYPE_AGGREGATE = 2; // 聚合
TRANSFORM_TYPE_JOIN = 3; // Join
TRANSFORM_TYPE_FILTER = 4; // 过滤
TRANSFORM_TYPE_UDF = 5; // 用户自定义函数
TRANSFORM_TYPE_EXPRESSION = 6; // 表达式计算
}
#8. 错误处理:gRPC Status 映射
适配器层将业务异常映射为标准的 gRPC Status:
| 业务异常 | gRPC Status | HTTP 等价 |
|---|---|---|
| PipelineAlreadyExistsException | ALREADY_EXISTS | 409 |
| IllegalArgumentException | INVALID_ARGUMENT | 400 |
| PipelineNotFoundException | NOT_FOUND | 404 |
| Exception(其他) | INTERNAL | 500 |
Metadata 容错:在构建响应时,对 Metadata 的每个字段都做了 null-safe 处理:
Timestamp createdAt = storedDef.hasMetadata()
? storedDef.getMetadata().getCreatedAt()
: Timestamp.getDefaultInstance();
使用 hasMetadata() 检查而非直接调用 getMetadata() 避免了 NullPointerException。
#9. Metadata 与版本管理
PipelineDefinition 的 Metadata 字段承载了丰富的管理信息:
message PipelineMetadata {
string created_by = 1;
google.protobuf.Timestamp created_at = 2;
string updated_by = 3;
google.protobuf.Timestamp updated_at = 4;
string version = 5;
map<string, string> labels = 6;
map<string, string> annotations = 7;
}
版本号的获取有优先级逻辑:
String versionId = storedDef.hasMetadata()
? storedDef.getMetadata().getVersion()
: storedDef.getVersion();
优先使用 Metadata 中的版本号(由服务端管理),回退到 PipelineDefinition 顶层的 version 字段(客户端指定)。
#10. Pipeline 监控与可观测性
Pipeline 执行提供了多维度的可观测性:
- Step 级别状态:每个 Step 独立追踪 PENDING/RUNNING/COMPLETED/FAILED/SKIPPED
- 执行时间线:每个 Step 记录 started_at 和 completed_at
- 资源使用:通过 Quarkus Micrometer 集成暴露指标
- 日志关联:通过 execution_id 串联整个 Pipeline 执行的所有日志
PipelineServiceGrpcImpl 在每个 RPC 入口都记录了结构化日志:
LOG.infof("CreatePipeline: request received");
LOG.infof("GetPipeline: id=%s", pipelineId);
#11. Key Takeaways
- 薄适配器模式:gRPC Impl 只负责参数校验、异常映射和 Proto 转换,业务逻辑完全在 Service 层。
- CDI 代理陷阱:Quarkus 中 Uni lambda 与 CDI 代理的交互需要特别注意——先同步调用再包装 Uni。
- Step-DAG 编译模型:通过
depends_on字段将 Step 列表编译为 DAG,创建时验证、运行时调度。 - 三层异常分类:ALREADY_EXISTS / INVALID_ARGUMENT / INTERNAL 覆盖了所有错误场景。
- Metadata 容错:
hasMetadata()+ 默认值回退的模式确保了即使 Metadata 缺失也不会崩溃。 - 自动血缘记录:Pipeline 与 LineageService 的集成使得数据追溯成为零成本操作。
#下一篇
S9-17:AuditService — 跨进程审计的 Kafka Consumer,我们将深入 Metadata & Governance Layer 的审计服务,了解决策追踪和合规性审计的实现。
Tags: #coomia-dip #source-code-reading #pipeline-service #dsl-to-dag #quarkus #grpc #Layer-f