Back to Blog

Pipeline DSL Design: Python Chaining API

Tags: #PipelineDSL #ChainAPI #ETL #Flink #TypeSafety #coomia-dip

CoomiaPublished on July 28, 20258 min read
Share this articleTwitter / X

Series: S3 Data Foundation · Article 18 | Level: Advanced | Reading Time: 20 min

Pipeline DSL Design: Python Chaining API

Tags: #PipelineDSL #ChainAPI #ETL #Flink #TypeSafety #coomia-dip

#TL;DR

The coomia-dip Pipeline DSL provides a Python chaining API that lets users define data processing pipelines declaratively. Through Pipeline.create().source().transform().sink().build() chain calls, users build complex ETL flows without writing low-level Flink/Spark code. This article fully dissects the Pipeline DSL's syntax design, type safety mechanism, compilation optimization, execution engine adaptation, and debugging toolchain across the complete implementation chain.

#1. Why Pipeline DSL

In the data processing landscape, users face two extremes:

Low-code/no-code tools: drag-and-drop UIs, easy but limited, complex logic hard to express, version control difficult.

Native programming APIs: Flink/Spark APIs directly, powerful but steep learning curve, verbose code, hard to integrate with platform.

Palantir Foundry's solution is Code Repositories + Transforms API — a Python DSL between the two. Users write concise Python functions, Foundry orchestrates execution.

The coomia-dip Pipeline DSL borrows this philosophy with these design goals:

  • Chain API: fluent Builder pattern, code-as-documentation
  • Type Safety: compile-time schema compatibility checks, not runtime errors
  • Engine-agnostic: same DSL code compiles to Flink Job, Spark Job, or DuckDB query
  • Visualizable: DSL-defined pipelines auto-generate DAG visualization
  • Version-controlled: Pipeline definitions are Python code, naturally Git-friendly
Code
+------------------------------------------------------------------+
|  Pipeline DSL Positioning                                         |
|                                                                   |
|  Low-code UI <--------- Pipeline DSL ---------> Native API        |
|  (easy/limited)        (balance point)          (powerful/complex) |
|                                                                   |
|  Features:                                                        |
|  - Python chain syntax, code-as-documentation                     |
|  - Schema type safety                                             |
|  - Multi-engine compile (Flink / Spark / DuckDB)                  |
|  - Auto DAG visualization                                         |
|  - Git version control friendly                                   |
+------------------------------------------------------------------+

#2. Core API Design

#2.1 Pipeline Chain API

Python
from onto_pipeline import Pipeline, Source, Transform, Sink
from onto_pipeline.transforms import filter, map, join, aggregate, window
from onto_pipeline.sources import mysql_cdc, kafka, csv_file
from onto_pipeline.sinks import iceberg, doris, ontology

# Define a complete Pipeline
pipeline = (
    Pipeline.create("equipment-sync")
    .description("Equipment data real-time sync pipeline")
    .owner("data-team")

    # Data source
    .source(
        mysql_cdc("factory_db.equipment")
        .host("mysql.internal")
        .port(3306)
        .startup_mode("initial")
    )

    # Transform chain
    .transform(
        filter(lambda row: row["status"] != "decommissioned")
    )
    .transform(
        map(lambda row: {
            **row,
            "display_name": f"{row['type']}-{row['serial_number']}",
            "is_active": row["status"] in ("running", "idle"),
        })
    )
    .transform(
        aggregate(
            group_by=["type", "region"],
            aggregations={
                "count": ("id", "count"),
                "avg_uptime": ("uptime_hours", "avg"),
            }
        )
    )

    # Data sink
    .sink(
        ontology("Device")
        .mapping({
            "id": "entity_id",
            "display_name": "name",
            "type": "device_type",
        })
        .on_conflict("UPDATE")
    )

    .build()
)

#2.2 Source Connectors

Python
class SourceConnectorRegistry:
    """Source connector registry"""

    CONNECTORS = {
        'mysql_cdc': MySQLCDCSource,
        'postgresql_cdc': PostgreSQLCDCSource,
        'kafka': KafkaSource,
        'csv_file': CSVFileSource,
        'json_file': JSONFileSource,
        'parquet_file': ParquetFileSource,
        'iceberg': IcebergSource,
        'api': RESTAPISource,
        'minio': MinIOSource,
    }

# MySQL CDC Source
source = (
    mysql_cdc("factory_db.equipment")
    .host("mysql.internal")
    .port(3306)
    .username("cdc_user")
    .password_from_secret("mysql-cdc-password")
    .startup_mode("initial")  # or "latest-offset"
    .include_schema_changes(True)
    .server_timezone("Asia/Shanghai")
)

# Kafka Source
source = (
    kafka("equipment-events")
    .bootstrap_servers("kafka:9092")
    .group_id("pipeline-equipment")
    .format("json")
    .schema_registry("http://schema-registry:8081")
    .starting_offset("earliest")
)

#2.3 Transform Operations

Python
# Filter
.transform(filter(lambda row: row["value"] > 0))
.transform(filter("status != 'deleted'"))  # SQL expression

# Map
.transform(map(lambda row: {
    **row,
    "full_name": f"{row['first_name']} {row['last_name']}",
    "processed_at": datetime.utcnow(),
}))

# Join
.transform(
    join(
        right=Pipeline.ref("department-lookup"),
        on="department_id = id",
        type="left"
    )
)

# Aggregate
.transform(
    aggregate(
        group_by=["region", "device_type"],
        aggregations={
            "total_count": ("id", "count"),
            "avg_uptime": ("uptime", "avg"),
            "max_temperature": ("temperature", "max"),
        }
    )
)

# Window
.transform(
    window(
        type="tumbling",
        size="5 minutes",
        aggregations={
            "event_count": ("id", "count"),
            "avg_value": ("value", "avg"),
        }
    )
)

# Flatten (nested JSON)
.transform(flatten("metadata", prefix="meta_"))

# Deduplicate
.transform(deduplicate(
    keys=["entity_id"],
    order_by="updated_at DESC",
    window="1 hour"
))

#3. Type Safety Mechanism

#3.1 Schema Propagation

Python
class SchemaTracker:
    """Track schema through the transform chain"""

    def __init__(self, source_schema: Schema):
        self._current_schema = source_schema

    def apply_transform(self, transform: Transform) -> Schema:
        if isinstance(transform, FilterTransform):
            # Filter doesn't change schema
            self._validate_filter_columns(
                transform.predicate, self._current_schema
            )
            return self._current_schema

        if isinstance(transform, MapTransform):
            new_schema = self._infer_map_output(
                transform.function, self._current_schema
            )
            self._current_schema = new_schema
            return new_schema

        if isinstance(transform, AggregateTransform):
            new_schema = Schema(
                columns={
                    **{col: self._current_schema.columns[col]
                       for col in transform.group_by},
                    **{alias: self._infer_agg_type(agg)
                       for alias, agg in transform.aggregations.items()},
                }
            )
            self._current_schema = new_schema
            return new_schema

        raise ValueError(f"Unknown transform: {type(transform)}")

#3.2 Compile-Time Validation

Python
class PipelineValidator:
    """Pipeline compile-time validator"""

    def validate(self, pipeline: Pipeline) -> list[ValidationError]:
        errors = []

        # 1. Source schema validation
        source_schema = pipeline.source.get_schema()
        if source_schema is None:
            errors.append(ValidationError(
                "Cannot infer source schema",
                stage="source"
            ))
            return errors

        # 2. Transform chain validation
        current_schema = source_schema
        for i, transform in enumerate(pipeline.transforms):
            try:
                current_schema = self._validate_transform(
                    transform, current_schema
                )
            except SchemaError as e:
                errors.append(ValidationError(
                    f"Transform {i} ({transform.name}): {e}",
                    stage=f"transform[{i}]"
                ))

        # 3. Sink compatibility validation
        sink_schema = pipeline.sink.get_expected_schema()
        if sink_schema:
            incompatible = self._check_schema_compatibility(
                current_schema, sink_schema
            )
            if incompatible:
                errors.append(ValidationError(
                    f"Output schema incompatible with sink: "
                    f"{incompatible}",
                    stage="sink"
                ))

        return errors

#4. Multi-Engine Compilation

#4.1 Engine Abstraction

Python
class PipelineCompiler(ABC):
    """Pipeline compiler base class"""

    @abstractmethod
    def compile(self, pipeline: Pipeline) -> ExecutableJob:
        ...

class FlinkCompiler(PipelineCompiler):
    """Compile Pipeline DSL to Flink Job"""

    def compile(self, pipeline: Pipeline) -> FlinkJob:
        env = StreamExecutionEnvironment.get_execution_environment()

        # Source
        source_stream = self._compile_source(pipeline.source, env)

        # Transforms
        current = source_stream
        for transform in pipeline.transforms:
            current = self._compile_transform(transform, current)

        # Sink
        self._compile_sink(pipeline.sink, current)

        return FlinkJob(env=env, name=pipeline.name)


class DuckDBCompiler(PipelineCompiler):
    """Compile Pipeline DSL to DuckDB SQL"""

    def compile(self, pipeline: Pipeline) -> DuckDBQuery:
        # Build SQL from transform chain
        sql = self._compile_source_sql(pipeline.source)

        for transform in pipeline.transforms:
            sql = self._wrap_transform_sql(transform, sql)

        sql = self._compile_sink_sql(pipeline.sink, sql)

        return DuckDBQuery(sql=sql, name=pipeline.name)

#4.2 Engine Selection

Code
Pipeline Engine Selection Strategy:

+------------------+--------------------+-------------------+
| Criteria          | Engine Selected     | Reason             |
+------------------+--------------------+-------------------+
| CDC source        | Flink              | Native CDC support |
| Kafka source      | Flink              | Stream processing  |
| File source <1GB  | DuckDB             | Fast local process |
| File source >1GB  | Spark              | Distributed        |
| Window aggregation| Flink              | Stream windowing   |
| Batch transform   | DuckDB / Spark     | Based on size      |
| Real-time sink    | Flink              | Low latency        |
+------------------+--------------------+-------------------+

#5. DAG Visualization

#5.1 Auto-Generated DAG

Python
class PipelineDAGVisualizer:
    """Pipeline DAG auto-visualizer"""

    def generate_dag(self, pipeline: Pipeline) -> DAGGraph:
        nodes = []
        edges = []

        # Source node
        source_node = DAGNode(
            id="source",
            label=pipeline.source.display_name,
            type="source",
            schema=pipeline.source.get_schema()
        )
        nodes.append(source_node)

        prev_id = "source"
        for i, transform in enumerate(pipeline.transforms):
            node = DAGNode(
                id=f"transform_{i}",
                label=transform.display_name,
                type="transform",
                schema=transform.output_schema
            )
            nodes.append(node)
            edges.append(DAGEdge(
                from_id=prev_id,
                to_id=node.id,
                row_estimate=transform.estimated_output_rows
            ))
            prev_id = node.id

        sink_node = DAGNode(
            id="sink",
            label=pipeline.sink.display_name,
            type="sink"
        )
        nodes.append(sink_node)
        edges.append(DAGEdge(from_id=prev_id, to_id=sink_node.id))

        return DAGGraph(nodes=nodes, edges=edges)
Code
DAG Visualization Example:

  [MySQL CDC: factory_db.equipment]
       |
       | (100K rows/day)
       v
  [Filter: status != 'decommissioned']
       |
       | (~90K rows/day)
       v
  [Map: add display_name, is_active]
       |
       | (~90K rows/day)
       v
  [Aggregate: GROUP BY type, region]
       |
       | (~50 groups)
       v
  [Ontology Sink: Device]

#6. Debugging and Monitoring

#6.1 Pipeline Dry Run

Python
class PipelineDryRunner:
    """Pipeline dry run — test with sample data"""

    def dry_run(
        self, pipeline: Pipeline, sample_size: int = 100
    ) -> DryRunResult:
        # Fetch sample from source
        sample = pipeline.source.sample(sample_size)

        results_per_stage = [("source", sample)]

        current = sample
        for i, transform in enumerate(pipeline.transforms):
            try:
                current = transform.apply(current)
                results_per_stage.append(
                    (f"transform_{i}: {transform.name}", current)
                )
            except Exception as e:
                return DryRunResult(
                    success=False,
                    failed_stage=i,
                    error=str(e),
                    stages=results_per_stage
                )

        return DryRunResult(
            success=True,
            stages=results_per_stage,
            output_schema=current.schema,
            output_sample=current.head(10)
        )

#6.2 Runtime Metrics

Code
Pipeline Runtime Metrics:

+------------------+----------+----------+----------+
| Metric            | Source   | Filter   | Sink     |
+------------------+----------+----------+----------+
| Records/sec       | 5,200    | 4,700    | 4,700    |
| Bytes/sec         | 2.1 MB   | 1.9 MB   | 1.5 MB   |
| Latency P50       | 12ms     | 2ms      | 45ms     |
| Latency P99       | 85ms     | 8ms      | 220ms    |
| Error rate         | 0.001%   | 0%       | 0.005%   |
| Backpressure       | 0%       | 0%       | 15%      |
+------------------+----------+----------+----------+

#7. Error Handling and Retry

Python
class PipelineErrorHandler:
    """Pipeline error handling"""

    def __init__(self, config: ErrorConfig):
        self._config = config

    def configure_pipeline(self, pipeline: Pipeline) -> Pipeline:
        return (
            pipeline
            .on_error(
                strategy="retry",
                max_retries=3,
                backoff="exponential",
                base_delay_ms=1000
            )
            .dead_letter_queue(
                sink=kafka("pipeline-dlq")
                .topic(f"dlq-{pipeline.name}")
            )
            .checkpoint(
                interval="1 minute",
                mode="exactly_once"
            )
            .metrics(
                enabled=True,
                export_to="prometheus"
            )
        )

#8. Testing Strategy

Python
class TestPipelineDSL:

    def test_simple_pipeline_builds(self):
        pipeline = (
            Pipeline.create("test")
            .source(csv_file("data.csv"))
            .transform(filter(lambda r: r["age"] > 0))
            .sink(doris("test_table"))
            .build()
        )
        assert pipeline.name == "test"
        assert len(pipeline.transforms) == 1

    def test_schema_validation_catches_errors(self):
        pipeline = (
            Pipeline.create("bad")
            .source(csv_file("data.csv").schema({"name": "str"}))
            .transform(filter(lambda r: r["age"] > 0))  # age not in schema
            .sink(doris("test_table"))
        )
        errors = pipeline.validate()
        assert len(errors) > 0
        assert "age" in errors[0].message

    def test_dry_run_with_sample(self):
        pipeline = (
            Pipeline.create("test")
            .source(csv_file("test_data.csv"))
            .transform(map(lambda r: {**r, "upper_name": r["name"].upper()}))
            .sink(doris("output"))
            .build()
        )
        result = pipeline.dry_run(sample_size=10)
        assert result.success
        assert "upper_name" in result.output_schema.columns

    def test_multi_engine_compilation(self):
        pipeline = Pipeline.create("test").source(...).transform(...).sink(...).build()

        flink_job = FlinkCompiler().compile(pipeline)
        assert flink_job is not None

        duckdb_query = DuckDBCompiler().compile(pipeline)
        assert duckdb_query.sql is not None

#Key Takeaways

  1. Pipeline DSL bridges the gap between low-code UIs and native APIs: a Python chaining syntax that is both readable as documentation and powerful enough for complex ETL logic.

  2. Type safety catches errors at compile time: schema propagation through the transform chain validates column references and type compatibility before any data flows.

  3. Engine-agnostic compilation maximizes flexibility: the same DSL definition can compile to Flink (streaming), Spark (batch), or DuckDB (embedded) based on workload characteristics.

  4. Dry run and DAG visualization accelerate development: test pipelines with sample data before deployment, and visualize the transform DAG to understand data flow.

  5. Built-in error handling and monitoring provide production-readiness: retry strategies, dead letter queues, checkpointing, and Prometheus metrics ensure reliable pipeline operation.

#Next Article

Next up: S3-19 "DolphinScheduler Integration: Workflow Orchestration Engine" will show how Pipeline DSL definitions are scheduled and orchestrated through DolphinScheduler for periodic batch and event-driven execution.

Tags: #PipelineDSL #ChainAPI #ETL #Flink #Spark #DuckDB #TypeSafety #DAGVisualization #StreamProcessing #coomia-dip #DataFoundation