Skip to content

Write an engine

Prerequisites · You have a DataFrame library you want to run DataCoolie pipelines on · you're ready to implement the full engine contract. End state · A new engine that implements the BaseEngine contract, passes its plugin-owned native tests and integration checks, and can be selected via create_engine("mylib").

Large surface area

Expect a substantial implementation and conformance effort. Start by studying datacoolie.engines.polars_engine.PolarsEngine as a behavioral template, then use the BaseEngine API reference as the authoritative abstract contract.

Skeleton

from datacoolie.engines.base import BaseEngine
from datacoolie.engines.contracts.windows import WindowSpec
import mylib


class MyLibEngine(BaseEngine[mylib.DataFrame]):
    def __init__(self, platform=None):
        super().__init__(platform)

    # --- Read ---
    def read_parquet(self, path, options=None): ...
    def read_delta(self, path, options=None):  ...
    def read_iceberg(self, path, options=None): ...
    # ... etc. for csv, json, jsonl, avro, excel

    def read_path(self, path, fmt, options=None):
        # Dispatch on fmt → the right abstract reader
        ...

    def read_database(self, *, table=None, query=None, options=None): ...
    def read_table(self, table_name, fmt="delta", options=None): ...
    def create_dataframe(self, records): ...
    def execute_sql(self, sql, parameters=None): ...

    # --- Write ---
    def write_to_path(self, df, path, mode, fmt, partition_columns=None, options=None): ...
    def write_to_table(self, df, table_name, mode, fmt, partition_columns=None, options=None): ...

    # --- Merge ---
    def merge_to_path(self, df, path, merge_keys, fmt="delta", partition_columns=None, options=None): ...
    def merge_overwrite_to_path(
        self, df, path, merge_keys, fmt="delta", partition_columns=None,
        options=None, write_options=None,
    ): ...
    def merge_to_table(self, df, table_name, merge_keys, fmt, partition_columns=None, options=None): ...
    def merge_overwrite_to_table(
        self, df, table_name, merge_keys, fmt="delta", partition_columns=None,
        options=None, write_options=None,
    ): ...

    # --- Transform, system columns, metrics, maintenance, SCD2 ---
    # (see BaseEngine for the full list)

The cast_column implementation receives the authored source declaration and its optional dialect context. It must resolve the declaration using the pure datacoolie.engines.data_types contract, then construct the plugin's native datatype or expression directly:

def cast_column(
    self, df, column_name, target_type, fmt=None, *,
    type_system=None, precision=None, scale=None,
): ...

Do not require SchemaConverter to pre-serialize a target string, and do not apply schema hints while reading or calculating watermarks. A plugin may use different native objects from Spark/Polars, but supported logical ranges, null/error behavior, and temporal semantics must remain equivalent.

The system-column contract includes the optional driver execution ID:

def add_system_columns(self, df, author=None, dataflow_run_id=None):
    # Add the standard timestamps/author. When dataflow_run_id is provided,
    # add it as the string column __dataflow_run_id.
    ...

Do not generate a replacement ID inside the engine: it must remain equal to the DataFlowRuntimeInfo.dataflow_run_id supplied by the driver.

fmt parameter contract

Every format-aware method must accept a fmt string. read_table, merge_to_table, and table_exists_by_name have contract-specific signatures:

def read_table(self, table_name: str, fmt: str = "delta", options=None): ...
def merge_to_table(
    self, df, table_name, merge_keys, fmt: str,
    partition_columns=None, options=None,
): ...
def table_exists_by_name(self, table_name: str, *, fmt: str = "delta") -> bool: ...

table_exists_by_name uses keyword-only fmt.

merge_overwrite_to_path and merge_overwrite_to_table receive write_options separately from options; preserve and forward both option maps to the native writer. Dropping write_options changes merge-overwrite behavior even when the merge keys are correct.

See ADR-0001.

See the BaseEngine API reference for the full abstract contract and navigation helpers.

Register

[project.entry-points."datacoolie.engines"]
mylib = "mypkg.engine:MyLibEngine"

Engine fmt support is a backend capability, while an entry point only adds a runtime registry name. These contracts are separate from authored metadata: the 0.2.0 schema enumerates built-in connection formats, and omitting connection_type is not a dc validate workaround for a custom format. Use a new entry-point alias; a packaged alias that collides with a built-in name does not replace the built-in registration.

Conformance

The framework's generic tests are contract references, not a substitute for a plugin-owned qualification suite. A new engine should maintain native unit tests for every abstract method it implements, including format dispatch, write_options forwarding, typed casts, null/error behavior, watermark filtering, metrics, and system columns. Add integration tests against the actual dataframe backend that register the engine, read and write at least one supported path/table format, exercise one supported merge or replacement operation, and verify the resulting rows and schema. Add explicit negative tests for unsupported formats and capabilities that prove they fail before mutating a target. Use the built-in Polars and Spark implementations as behavioral references and run only the framework tests applicable to the contract surface your plugin claims.

replace_window — Engine-owned bounded replacement

Engines may override the concrete replace_window operation to provide a native merge/transaction or format-specific materialization. The base implementation validates the input window, deletes rows in the bounded scope, and appends the non-empty input. It intentionally does not claim atomicity.

def replace_window(
    self, df, *, table_name=None, path=None, window, fmt="delta",
    partition_columns=None, options=None,
):
    ...

The neutral WindowSpec carries bound values, lower/upper operators, and OR composition across columns. Engines render identifiers and literals for their own backend. Validate the final transformed watermark columns before deleting; do not resolve schema hints or infer arbitrary custom-transform mappings here.

delete_by_window — Range-based delete primitive

Engines must implement delete_by_window_path and delete_by_window_table to support the replace_by_watermark destination feature:

@abstractmethod
def delete_by_window_path(
    self,
    path: str,
    window: WindowSpec,
    fmt: str = "delta",
) -> None:
    """Delete rows in a path-based table within the value window."""

@abstractmethod
def delete_by_window_table(
    self,
    table_name: str,
    window: WindowSpec,
    fmt: str = "delta",
) -> None:
    """Delete rows in a named table within the value window."""

WindowSpec.bounds maps column names to (lower_bound, upper_bound) tuples. Build a predicate like col > lower AND col <= upper for each entry and delete all matching rows. The lower/upper operators come from the WindowSpec and must be preserved; the default lower bound is exclusive to match the source watermark filter semantics. Predicates are OR-combined across columns and AND-combined within one column.

The engine must materialize or checkpoint a lazy input before deletion when a later append would otherwise re-evaluate it against the changed target. Release only resources created by that operation, and do not hide the original write error if cleanup also fails.

The base replacement method may use these primitives, and a backend can route to them internally. New destination strategies should call replace_window, not sequence public delete and append operations themselves.