返回博客

源码精读: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 编译模型。

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

源码精读: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 编译模型。

#目录

  1. 整体架构:Pipeline & Orchestration Layer 的 Quarkus gRPC 服务
  2. Proto 合约:PipelineDefinition 与 Step-DAG 模型
  3. PipelineServiceGrpcImpl:适配器模式
  4. Quarkus CDI 代理的作用域陷阱
  5. CreatePipeline:定义与验证
  6. Pipeline 执行:DAG 调度策略
  7. 血缘记录:Pipeline 与 LineageService 的集成
  8. 错误处理:gRPC Status 映射
  9. Metadata 与版本管理
  10. Pipeline 监控与可观测性
  11. Key Takeaways

#1. 整体架构:Pipeline & Orchestration Layer 的 Quarkus gRPC 服务

Pipeline 服务位于 Data Layer(Data Layer + Pipeline & Orchestration Layer 合并后),使用 Quarkus 3.x 框架构建:

Code
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 组成:

PROTOBUF
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 是一个典型的薄适配器,不包含业务逻辑:

Java
@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 的陷阱:

Java
/**
 * 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 创建流程包含多层验证:

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

三层异常处理形成了清晰的错误分类:

  • PipelineAlreadyExistsExceptionALREADY_EXISTS:Pipeline 名称唯一性冲突
  • IllegalArgumentExceptionINVALID_ARGUMENT:输入验证失败(如 DAG 包含环)
  • ExceptionINTERNAL:未预期的内部错误

#6. Pipeline 执行:DAG 调度策略

Pipeline 执行使用拓扑排序确定 Step 的执行顺序:

  1. 拓扑排序:解析 depends_on 构建 DAG,进行拓扑排序
  2. 并行调度:同一拓扑层级(无依赖关系)的 Step 可并行执行
  3. 失败处理:某个 Step 失败时,其所有下游 Step 自动标记为 SKIPPED
  4. 重试机制:支持 Step 级别的重试策略

DAG 的编译发生在 Pipeline 创建时——PipelineExecutionService.createPipeline() 会验证 DAG 的合法性(无环、所有引用的 step_id 存在)。运行时只需按照拓扑序调度。

#7. 血缘记录:Pipeline 与 LineageService 的集成

Pipeline 执行过程中会自动记录血缘信息到 Pipeline & Orchestration Layer 的 LineageService:

PROTOBUF
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 支持六种转换类型:

PROTOBUF
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 StatusHTTP 等价
PipelineAlreadyExistsExceptionALREADY_EXISTS409
IllegalArgumentExceptionINVALID_ARGUMENT400
PipelineNotFoundExceptionNOT_FOUND404
Exception(其他)INTERNAL500

Metadata 容错:在构建响应时,对 Metadata 的每个字段都做了 null-safe 处理:

Java
Timestamp createdAt = storedDef.hasMetadata()
    ? storedDef.getMetadata().getCreatedAt()
    : Timestamp.getDefaultInstance();

使用 hasMetadata() 检查而非直接调用 getMetadata() 避免了 NullPointerException。

#9. Metadata 与版本管理

PipelineDefinition 的 Metadata 字段承载了丰富的管理信息:

PROTOBUF
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;
}

版本号的获取有优先级逻辑:

Java
String versionId = storedDef.hasMetadata()
    ? storedDef.getMetadata().getVersion()
    : storedDef.getVersion();

优先使用 Metadata 中的版本号(由服务端管理),回退到 PipelineDefinition 顶层的 version 字段(客户端指定)。

#10. Pipeline 监控与可观测性

Pipeline 执行提供了多维度的可观测性:

  1. Step 级别状态:每个 Step 独立追踪 PENDING/RUNNING/COMPLETED/FAILED/SKIPPED
  2. 执行时间线:每个 Step 记录 started_at 和 completed_at
  3. 资源使用:通过 Quarkus Micrometer 集成暴露指标
  4. 日志关联:通过 execution_id 串联整个 Pipeline 执行的所有日志

PipelineServiceGrpcImpl 在每个 RPC 入口都记录了结构化日志:

Java
LOG.infof("CreatePipeline: request received");
LOG.infof("GetPipeline: id=%s", pipelineId);

#11. Key Takeaways

  1. 薄适配器模式:gRPC Impl 只负责参数校验、异常映射和 Proto 转换,业务逻辑完全在 Service 层。
  2. CDI 代理陷阱:Quarkus 中 Uni lambda 与 CDI 代理的交互需要特别注意——先同步调用再包装 Uni。
  3. Step-DAG 编译模型:通过 depends_on 字段将 Step 列表编译为 DAG,创建时验证、运行时调度。
  4. 三层异常分类:ALREADY_EXISTS / INVALID_ARGUMENT / INTERNAL 覆盖了所有错误场景。
  5. Metadata 容错hasMetadata() + 默认值回退的模式确保了即使 Metadata 缺失也不会崩溃。
  6. 自动血缘记录: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