Pipeline DSL Design: Python Chaining API
Tags: #PipelineDSL #ChainAPI #ETL #Flink #TypeSafety #coomia-dip
“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
+------------------------------------------------------------------+
| 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
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
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
# 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
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
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
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
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
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)
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
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
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
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
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
-
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.
-
Type safety catches errors at compile time: schema propagation through the transform chain validates column references and type compatibility before any data flows.
-
Engine-agnostic compilation maximizes flexibility: the same DSL definition can compile to Flink (streaming), Spark (batch), or DuckDB (embedded) based on workload characteristics.
-
Dry run and DAG visualization accelerate development: test pipelines with sample data before deployment, and visualize the transform DAG to understand data flow.
-
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