Write a transformer¶
Prerequisites · You have a per-row or per-batch transformation to apply between read and write. End state · Transformer in the pipeline at a known order slot, registered via entry points.
For a complete local package installation and Driver run, start with the transformer tutorial. This page explains the implementation contract and ordering choices. For ordinary masking, first consider the built-in transform patterns.
Minimal transformer¶
The complete source for this pattern is available as source or raw. Its packaging metadata is raw.
from datacoolie.transformers.base import BaseTransformer
from datacoolie.core.models.dataflow import DataFlow
class PiiMaskerTransformer(BaseTransformer):
# Slots 40-50 are reserved for user plugins. Pick one.
ORDER = 45
def __init__(self, engine) -> None:
self._engine = engine
@property
def order(self) -> int:
return self.ORDER
def transform(self, df, dataflow: DataFlow):
cfg = dataflow.transform.configure.get("pii_mask", {})
cols = cfg.get("columns", [])
if not cols:
self._mark_skipped()
return df
for c in cols:
df = self._engine.add_column(
df, c, f"CASE WHEN {c} IS NULL THEN NULL ELSE '***' END"
)
self._mark_applied(f"cols={len(cols)}")
return df
Register¶
This matches the downloadable flat module pii_masker.py. If your package uses
another module path, change the entry-point value to its importable module and
class. Install the package in the Python environment that runs the Driver.
Opt in from metadata¶
Registering a transformer makes it resolvable, but does not add it to the framework's default transformer sequence. Use the protected Driver hook shown below; check its API signature and rerun your integration test when upgrading:
from datacoolie import transformer_registry
from datacoolie.core.constants import ColumnCaseMode
from datacoolie.orchestration.driver import DataCoolieDriver
class PiiDriver(DataCoolieDriver):
def _create_transformer_pipeline(
self,
dataflow_run_id=None,
column_name_mode=ColumnCaseMode.LOWER,
):
pipeline = super()._create_transformer_pipeline(
dataflow_run_id=dataflow_run_id,
column_name_mode=column_name_mode,
)
pipeline.add_transformer(
transformer_registry.get("pii_masker", engine=self._engine)
)
return pipeline
driver = PiiDriver(engine=engine, metadata_provider=metadata)
Here engine and metadata are already configured instances, and the engine
has a platform attached. Otherwise supply platform= to the Driver; both
must share the same platform object. The
tutorial supplies a complete
runner. Its metadata uses transform.configure.pii_mask.columns; that map is
read by this plugin after the runner adds it to the pipeline.
TransformerPipeline sorts the combined list by each transformer's order
when it runs.
Always forward both dataflow_run_id and column_name_mode when overriding
this hook. The driver uses the run ID to correlate the framework-owned
__dataflow_run_id column with the DataFlowRuntimeInfo recorded for the same
execution, and the mode controls column-name normalization.
Order slot cheat-sheet¶
| Slots | Who owns them |
|---|---|
| 5 | ColumnValueTransformer |
| 10 | SchemaConverter |
| 18 | HashColumnAdder |
| 20 | Deduplicator |
| 30 | ColumnAdder |
| 35 | RowFilter |
| 40–50 | Your plugins |
| 60 | SCD2ColumnAdder |
| 70 | SystemColumnAdder |
| 80 | PartitionHandler |
| 84 | DataMasker |
| 85 | ColumnProjector |
| 90 | ColumnNameSanitizer |
| Other slots | Reserved for future framework work or additional plugins; the pipeline sorts by numeric order |
See ADR-0003.
Tracking labels¶
Call _mark_applied(), _mark_applied("detail"), or _mark_skipped() inside
transform so the execution log records exactly what your transformer did. Without
a call, the default is to record your class name.
Test your implementation¶
Test with native DataFrames for every advertised engine. Check configured columns, preserved nulls, unchanged output when no columns are configured, ordering with other transforms, and failure for invalid column expressions. The sample assumes simple SQL-safe column names; adapt expression construction to the identifiers and data types your plugin supports. Do not treat it as a general identifier-quoting or anonymization implementation.
Then test an installed package through a Driver run, using the same alias as its entry point and checking business output. See the discovery troubleshooting steps and transformer API.