Transform patterns¶
Prerequisites · A configured source from Source patterns.
Choose the destination load strategy before finalizing transforms when the case
uses merge keys, SCD2 columns, partitions, or watermark replacement.
End state · A correct transform block in each dataflow that needs data shaping.
Each section below follows one transformer class: what it does, which metadata configures it, and an example. The table links to those sections in execution order. For source-aware datatype rules, use Datatypes and schema hints.
The Transform block is optional.
Most business transformations are configured there. SCD2 and partition columns
use destination metadata; system columns are automatic, and column-name
sanitization uses a Driver run option. The built-in pipeline runs between read
and write in this order:
| Order | Transformer class | Metadata / configuration |
|---|---|---|
| 5 | ColumnValueTransformer |
transform.value_rules |
| 10 | SchemaConverter |
transform.schema_hints; timestamp conversion in transform.configure |
| 18 | HashColumnAdder |
transform.hash_columns |
| 20 | Deduplicator |
transform.deduplicate_columns, latest_data_columns; merge/watermark fallbacks |
| 30 | ColumnAdder |
transform.additional_columns; automatic stale-system-column cleanup |
| 35 | RowFilter |
transform.filter_expression |
| 60 | SCD2ColumnAdder |
destination.load_type = "scd2" and destination.configure.scd2_effective_column |
| 70 | SystemColumnAdder |
Automatic |
| 80 | PartitionHandler |
destination.partition_columns |
| 84 | DataMasker |
transform.masking_rules |
| 85 | ColumnProjector |
transform.select_columns, drop_columns, rename_columns |
| 90 | ColumnNameSanitizer |
Automatic; Driver column_name_mode |
System columns are always added
__created_at, __updated_at, __updated_by, and
__dataflow_run_id are added to every driver-managed dataflow output.
You do not configure them — just expect them in the destination table.
ColumnValueTransformer¶
When to use: incoming strings, nulls, or category codes need cleanup before
datatype conversion. Configure transform.value_rules. Each
Value rule
specifies an operation, target columns, and any operation-specific options.
The rule's own order controls ordering within value_rules; the whole
value-rule transformer still runs at pipeline order 5.
The six operations below are individual items in transform.value_rules.
All except fill_null require string columns. Existing nulls stay null unless
you explicitly use fill_null.
trim: remove surrounding spaces¶
" Alice " becomes "Alice"; " " becomes "". Only ASCII U+0020
spaces are removed. Tabs, newlines, and non-breaking spaces remain. One rule
can target several columns.
case: lowercase or uppercase¶
Use lower for normalized emails or upper for codes:
"transform": {
"value_rules": [
{"operation": "case", "columns": ["email"], "mode": "lower"},
{"operation": "case", "columns": ["country_code"], "mode": "upper"}
]
}
"Alice@Example.COM" becomes "alice@example.com"; "vn" becomes "VN".
These are the two supported modes; case conversion does not trim spaces.
regex_replace: replace every matching substring¶
Remove formatting from a phone number:
"(+84) 123-456" becomes "84123456". Set a non-empty replacement to
replace matches with text; omitting it uses the empty string.
The pattern uses DataCoolie portable regex v1 rather than the complete Java
or Rust dialect. Use literals, explicit character classes/ranges, ., anchors,
grouping, alternation, and ordinary quantifiers. Lookaround, backreferences,
named groups, inline flags, \d/\w/\s/\b, possessive quantifiers, and
quantified nested groups are rejected while loading metadata. Patterns are
limited to 4,096 characters. Replacement text is always literal, so $ and
backslash are emitted as written rather than expanding capture groups.
empty_to_null: convert empty strings to null¶
Only "" becomes null. " " remains unchanged unless an earlier trim
rule first removes its spaces. Non-empty values remain unchanged.
fill_null: supply a typed default¶
"transform": {
"value_rules": [
{"operation": "fill_null", "columns": ["country_code"], "value": "UNKNOWN"},
{"operation": "fill_null", "columns": ["quantity"], "value": 0},
{"operation": "fill_null", "columns": ["is_active"], "value": false}
]
}
Only nulls change. The example assumes country_code is already a string,
quantity an integer, and is_active a boolean. value must match the
current column type; schema hints run later and cannot make an incompatible
literal valid at this stage.
| Current column type | Example JSON value |
Requirement |
|---|---|---|
| String | "UNKNOWN" |
Use a JSON string, including "" if intentional |
| Integer | 0 |
Integer within the target range; "0" and booleans are rejected |
| Float | 0.5 |
Finite JSON number |
| Decimal | "0.00" |
Integer or decimal string that fits the column's precision/scale |
| Boolean | false |
JSON boolean, not "false" or 0 |
| Date | "2026-01-01" |
ISO date |
| Timezone-free timestamp | "2026-01-01T00:00:00" |
ISO datetime without an offset |
| Timezone-aware timestamp | "2026-01-01T00:00:00+00:00" |
ISO datetime with an offset |
Binary, nested, and untyped-null columns are not supported literal targets. Invalid literals fail before the native expression is applied, independently of Spark ANSI settings. Empty strings and NaN are not nulls for this operation.
map: translate string codes¶
{"operation": "map", "columns": ["status"], "mapping": {"A": "active", "I": "inactive"}, "on_unmapped": "keep"}
"A" becomes "active"; an unknown "X" stays "X" with keep, the
default. Keys and values must be strings, and matching is exact: "a" does
not match "A". Use an earlier case rule if needed.
Use null when unknown codes should become null:
{
"operation": "map",
"columns": ["status"],
"mapping": {"A": "active", "I": "inactive"},
"on_unmapped": "null"
}
Only keep and null are supported. This setting controls unmapped values,
not absent columns.
Combine rules in a deliberate order¶
The following turns spaces-only input into "UNKNOWN":
"transform": {
"value_rules": [
{"operation": "trim", "columns": ["country_code"], "order": 10},
{"operation": "empty_to_null", "columns": ["country_code"], "order": 20},
{"operation": "fill_null", "columns": ["country_code"], "value": "UNKNOWN", "order": 30}
],
"configure": {"missing_column_policy": "error"}
}
Rules default to order 100; ties retain declaration order. An empty or omitted
value_rules list does nothing. Missing configured columns follow
transform.configure;
see Missing-column policy for the error/ignore cases.
SchemaConverter¶
When to use: your source data has weak types (CSV strings, JSON mixed types) and you need specific types in the destination.
Cast selected columns in one dataflow¶
Add schema_hints inside transform. The Schema hint
entry defines each item:
"transform": {
"schema_hints": [
{ "column_name": "order_id", "data_type": "int" },
{ "column_name": "customer_id", "data_type": "int" },
{ "column_name": "amount", "data_type": "decimal(18,2)" },
{ "column_name": "order_date", "data_type": "date" },
{ "column_name": "created_at", "data_type": "timestamp" },
{ "column_name": "is_active", "data_type": "boolean" }
]
}
For source-specific type names, decimals, timezone conversion, and the choice
between global and dataflow hints, use Datatypes and schema hints.
Reusable root-level schema_hints groups follow the
Shared schema hint
contract and are attached to matching dataflows by the metadata provider.
Less-common hint fields such as format, default_value, ordinal_position,
and is_active are listed in the Schema hint reference.
How schema hints are applied
- Matching is case-insensitive
- Missing columns are skipped, not fatal
use_schema_hintin the source connection configure must be enabled for hint-based caststimestamp_ntzconversion happens after hint-based casting when enabled withtimestamp_timezone; see Timestamp semantics
Only list the columns you want to cast
You do not need a hint for every column. Columns not listed keep their inferred type from the source reader.
Reuse global hints or override them locally¶
At the metadata document root, a shared group can describe the source table:
"schema_hints": [
{
"connection_name": "erp_postgres",
"schema_name": "sales",
"table_name": "orders",
"hints": [
{"column_name": "order_id", "data_type": "bigint"},
{"column_name": "amount", "data_type": "numeric(18,2)"}
]
}
]
The provider attaches this group to a source using that connection, schema,
and table. A non-empty transform.schema_hints list replaces the whole shared
group for that dataflow. To keep some shared casts, repeat them in the local
list. See Global hints for a source table
for matching rules, including query/function sources.
Decimal precision and formatted dates¶
"transform": {
"schema_hints": [
{"column_name": "amount", "data_type": "decimal", "precision": 18, "scale": 2},
{"column_name": "order_date", "data_type": "date", "format": "dd/MM/yyyy"},
{"column_name": "created_at", "data_type": "timestamp_ntz", "format": "dd/MM/yyyy HH:mm:ss"}
]
}
decimal with precision: 18, scale: 2 is equivalent to decimal(18,2).
If both forms supply parameters, they must agree. The format examples parse
"25/09/2026" and "25/09/2026 11:10:00". Use formats supported by the
selected engine; these common Spark-style patterns are translated by Polars.
Arbitrary vendor format strings are not a portable contract.
Interpret hints from a source datatype system¶
When a weak source reuses database type declarations, set the source connection configure:
For example, a PostgreSQL int8 hint then means a signed 64-bit integer.
Database connections normally supply their dialect through database_type.
See Which type system is read?
for all supported systems and precedence.
Disable one hint or all hint-based casts¶
Keep a hint in metadata while temporarily leaving its column unchanged:
"transform": {
"schema_hints": [
{"column_name": "amount", "data_type": "decimal(18,2)", "is_active": false}
]
}
To disable both global and inline casts for a source connection, set its
configure as follows:
Missing hinted columns are warned about and skipped; duplicate hint names are
rejected case-insensitively. default_value does not fill missing/null data in
the built-in converter; use ColumnValueTransformer.fill_null for that.
ordinal_position orders shared hint records during provider resolution; it
does not reorder output columns. Use ColumnProjector.select_columns to do so.
Convert timezone-free timestamps¶
To interpret NTZ values as instants, set convert_timestamp_ntz and
timestamp_timezone in transform.configure:
"transform": {
"configure": {
"convert_timestamp_ntz": true,
"timestamp_timezone": "Asia/Ho_Chi_Minh"
}
}
This conversion runs after hint-based casts and can also apply to NTZ columns already present in the source schema. The timezone identifies the source wall-clock time; see Timestamp semantics for supported timezone forms and a conversion example.
To retain timezone-free values, leave the default or set:
Hint-based casting and NTZ conversion are separate switches: disabling
use_schema_hint does not disable an explicitly enabled NTZ conversion.
Columns already carrying timezone-aware timestamps retain their instant.
HashColumnAdder¶
When to use: you need a stable row or business-key identifier from existing
columns. Hashing runs after schema conversion but before deduplication and
computed columns. Its inputs can use normalized, cast source columns, but not
columns created later by additional_columns.
Each transform.hash_columns item follows the Hash column
shape: declare a target_column, ordered input columns, and an algorithm.
Hashing does not infer inputs from deduplicate_columns or merge_keys.
Targets must be new columns. Missing inputs follow missing_column_policy in
transform.configure.
Business-key hashes: SHA-256 or XXHash64¶
"transform": {
"hash_columns": [
{
"target_column": "customer_hash",
"columns": ["country_code", "customer_id"],
"algorithm": "sha256"
},
{
"target_column": "customer_key",
"columns": ["country_code", "customer_id"],
"algorithm": "xxhash64"
}
]
}
Choose the hash use case deliberately¶
| Use case | Input columns | Choice and caution |
|---|---|---|
| Compact surrogate-key-style value | Complete natural/business key | xxhash64 is compact, but can collide; use an identity/mapping table for authoritative uniqueness. |
| Cross-system identifier | Explicitly ordered stable identifier columns | Prefer sha256 when collision resistance matters. Keep column order and types consistent across producers. |
| Row hash / hashdiff | Non-key attributes whose changes matter | Prefer sha256; omit volatile audit or ingestion timestamps unless they define a change. |
| Low-entropy PII protection | Do not use an unkeyed hash | Plain SHA-256 is vulnerable to dictionary attack; use masking or an approved keyed pseudonymization service. |
| File/payload integrity | Not a row-hash use case | Use a digest of object bytes, not hash_columns. |
For a surrogate-key-style hash, explicitly repeat the intended business key:
{
"destination": {
"load_type": "merge_upsert",
"merge_keys": ["country_code", "customer_id"]
},
"transform": {
"deduplicate_columns": ["country_code", "customer_id"],
"hash_columns": [{
"target_column": "customer_sk",
"columns": ["country_code", "customer_id"],
"algorithm": "xxhash64"
}]
}
}
This repetition prevents a later deduplication or merge change from silently changing persisted hash values. A generated hash target can itself be used as a merge key, but never as its own hash input. If deduplication columns and merge keys differ, choose the hash input deliberately.
Row hash / hashdiff for change detection¶
Hash the attributes whose changes matter separately from the identity:
"transform": {
"hash_columns": [
{"target_column": "customer_key", "columns": ["country_code", "customer_id"], "algorithm": "xxhash64"},
{"target_column": "customer_hashdiff", "columns": ["customer_name", "is_active", "signup_date"], "algorithm": "sha256"}
]
}
Here the hashdiff inputs are string, boolean, and Date columns. The generated value does not automatically enable CDC or change the destination merge/SCD2 strategy; configure its use separately. Multiple definitions run in declaration order, and all target names must be distinct and non-reserved.
Use a generated hash for deduplication and merging¶
Hashing runs before deduplication, so the new column can be the key:
{
"destination": {"load_type": "merge_upsert", "merge_keys": ["customer_hash"]},
"transform": {
"hash_columns": [{"target_column": "customer_hash", "columns": ["country_code", "customer_id"]}],
"deduplicate_columns": ["customer_hash"],
"latest_data_columns": ["updated_at"]
}
}
Omitting algorithm selects sha256. Keep this generated key in any later
projection and do not mask or rename it while it is a merge key.
Input types, serialization, and migration¶
Hash inputs currently support string, integer, boolean, and date columns. Input
order matters. The canonical payload uses type tags, null markers, and UTF-8
byte lengths, so Spark and Polars distinguish null from an empty string and
produce identical output. sha256 returns a lowercase 64-character String.
xxhash64 uses fixed seed 42 and returns a signed BIGINT; negative values are
normal. Only dc_hash_v1 serialization is currently supported; seed and salt
are not configurable metadata fields. Null input is encoded, so the hash
itself is still a value; null and empty string produce different payloads.
Polars loads polars-hash only when this
feature runs; install datacoolie[polars-hash] if needed.
XXHash64 is not collision-free. Do not apply abs() or discard its sign bit,
which reduces the key space. Add a collision quality check when using it as a
key. Changing an existing target from SHA-256 to XXHash64 also changes its type
from String to BIGINT; use a new target or plan a destination schema migration.
Plain SHA-256 is not protection for low-entropy PII.
Decimal, float, timestamp, binary, and nested inputs are not supported directly.
If those values must participate, define their supported representation
deliberately in the source query/function or a suitable earlier schema cast.
additional_columns runs too late to prepare a hash input. An empty or omitted
hash_columns list adds nothing.
Deduplicator¶
When to use: your source can deliver duplicate rows for the same key (common with CDC feeds, API pagination overlaps, or file re-deliveries).
Set deduplicate_columns and latest_data_columns in the
Transform block:
| Field | Meaning |
|---|---|
deduplicate_columns |
The column(s) that define a unique record — usually your natural key |
latest_data_columns |
Which column to use to pick the "winner" when duplicates exist — usually a timestamp |
Deduplicator groups by deduplicate_columns, orders by latest_data_columns
descending, and keeps the first row per group.
Composite keys and deterministic tie-breaking¶
"transform": {
"deduplicate_columns": ["country_code", "customer_id"],
"latest_data_columns": ["updated_at", "event_sequence"]
}
The ordering is descending across the tuple: newest updated_at, then highest
event_sequence. Both ordering columns must already exist. If all ordering
values tie, the default single winner is not a deterministic business choice;
add a tie-breaker or explicitly keep ties with rank.
Fallback behavior you should know¶
If you omit some dedup fields, DataCoolie still has a few convenience fallbacks:
| Missing input | Fallback |
|---|---|
deduplicate_columns |
Falls back to destination merge_keys |
latest_data_columns |
Falls back to source.watermark_columns |
| Either effective column list is empty after fallbacks | Deduplication becomes a no-op |
The fallback fields are described under Destination for merge keys and Source for watermarks.
Full example with merge_upsert:
{
"name": "orders_cdc_to_bronze",
"stage": "ingest",
"source": {
"connection_name": "cdc_source",
"table": "orders_changes",
"watermark_columns": ["updated_at"]
},
"destination": {
"connection_name": "bronze",
"schema_name": "sales",
"table": "orders",
"load_type": "merge_upsert",
"merge_keys": ["order_id"]
},
"transform": {
"deduplicate_columns": ["order_id"],
"latest_data_columns": ["updated_at"],
"schema_hints": [
{ "column_name": "order_id", "data_type": "long" },
{ "column_name": "updated_at", "data_type": "timestamp" }
]
}
}
Relationship to merge_keys
deduplicate_columns is usually the same value as merge_keys but it
lives in transform, not destination. They serve different pipeline
stages: deduplication happens before the merge.
Keep ties with rank instead of row-number¶
Normally DataCoolie keeps a single winner per key. If you want rank-style
deduplication instead, enable deduplicate_by_rank in
transform.configure:
"transform": {
"deduplicate_columns": ["order_id"],
"latest_data_columns": ["updated_at"],
"configure": {
"deduplicate_by_rank": true
}
}
merge_overwrite also uses rank-based dedup automatically when merge keys are
available and explicit deduplicate_columns are not set.
For example, this fragment takes its grouping and ordering from destination and source metadata and keeps tied latest rows:
{
"source": {"watermark_columns": ["updated_at"]},
"destination": {"load_type": "merge_overwrite", "merge_keys": ["order_id"]},
"transform": {}
}
For a key with ordering values 10, 10, 9, rank keeps both rows at 10;
row-number keeps one. To use row-number for this merge_overwrite case, supply
explicit deduplicate_columns and leave deduplicate_by_rank false. Setting
deduplicate_by_rank: true explicitly keeps ties for other load types too.
With multiple ordering columns, ties are compared across the full tuple.
An empty transform does not disable deduplication when both fallback lists
exist. If either effective list is empty it skips; if a configured column is
absent, it fails even when missing_column_policy is ignore.
ColumnAdder¶
When to use: you need a new column whose value is calculated from existing columns (derived date parts, string concatenation, status labels, etc.).
Each transform.additional_columns item follows the Additional column
shape: it defines a column
and its expression.
"transform": {
"additional_columns": [
{ "column": "order_year", "expression": "EXTRACT(YEAR FROM order_date)" },
{ "column": "order_month", "expression": "EXTRACT(MONTH FROM order_date)" },
{ "column": "full_name", "expression": "first_name || ' ' || last_name" },
{ "column": "is_large", "expression": "CASE WHEN amount > 1000 THEN true ELSE false END" }
]
}
Expressions are SQL evaluated against the DataFrame after schema casting.
Use standard SQL scalar functions — the Polars and Spark engines both support
EXTRACT, CASE WHEN, string functions, and arithmetic.
The class also removes stale framework system columns from the incoming data,
even when additional_columns is empty. SystemColumnAdder adds the current
execution's audit columns later.
Constants, replacing a column, and dependent expressions¶
Items run in declaration order. You can add a constant, replace an existing business column, then use that result in the next expression:
"transform": {
"additional_columns": [
{"column": "source_system", "expression": "'ERP'"},
{"column": "amount", "expression": "amount * 100"},
{"column": "is_large", "expression": "CASE WHEN amount > 100000 THEN true ELSE false END"}
]
}
The last expression sees the updated amount. SQL string constants need SQL
quotes inside the JSON string. Each expression must be non-empty. Use scalar
SQL expressions supported by the selected engine; arbitrary queries, joins,
and Python function calls belong in the source query/function workflow.
Polars SQL limitations
Polars does not support current_timestamp() or NOW().
Use EXTRACT(YEAR FROM col) instead of year(col).
Use CAST(col AS DATE) instead of date(col).
Do not reference system columns here
additional_columns runs at transformer order 30. System columns are only
added later at order 70, so expressions here cannot rely on __created_at,
__updated_at, __updated_by, or __dataflow_run_id. Let the framework add
those columns for you and use them after the transform stage, not inside it.
RowFilter¶
Configure transform.filter_expression
when you need to drop rows before writing to the destination.
The Transform reference
defines this field.
The predicate runs at order 35, after ColumnAdder (30), so it can
reference columns you created in additional_columns.
filter_expression is a SQL predicate — everything you would normally
place after WHERE. Both Polars and Spark engines evaluate it against the
DataFrame.
Reference a computed column¶
Because RowFilter runs after ColumnAdder, you can filter on columns
added by additional_columns:
"transform": {
"additional_columns": [
{ "column": "order_year", "expression": "EXTRACT(YEAR FROM order_date)" }
],
"filter_expression": "order_year >= 2023"
}
Combine multiple conditions¶
"transform": {
"filter_expression": "region = 'US' AND status NOT IN ('cancelled', 'draft') AND amount > 0"
}
vs source.filter_expression
There are two distinct filter hooks:
| Field | Stage | Scope |
|---|---|---|
source.filter_expression |
Read time (combined with or after the logical watermark condition) | Reader output columns, including aliases returned by source.query; not columns added by transforms |
transform.filter_expression |
Order 35 (post-ColumnAdder) | Source columns + computed columns |
Use source.filter_expression when you want the filter pushed as close to
the source as possible. Use transform.filter_expression when you need
to filter on a column that is added by additional_columns.
Null conditions and skipping the filter¶
Only rows for which the predicate is true survive. Use IS NULL / IS NOT NULL
for null checks; comparisons such as amount >= 0 do not retain null amounts.
Omit the field or use null to leave the row set unchanged:
Filtering runs before the new SCD2 and system columns are added, so its expression cannot depend on those columns.
SCD2ColumnAdder¶
When to use: the destination keeps Type 2 history. Set load_type to
scd2 in Destination, and set
scd2_effective_column in
destination.configure:
"destination": {
"connection_name": "gold",
"table": "customer",
"load_type": "scd2",
"merge_keys": ["customer_id"],
"configure": {"scd2_effective_column": "updated_at"}
}
At order 60, this class adds __valid_from from the effective column,
__valid_to as null, and __is_current as true for the destination writer.
The effective column must be available at this stage. No transform setting
is needed. See SCD2 configuration
and incremental SCD2 for the history-writing workflow.
For an incoming row with updated_at = 2026-09-25 11:10:00, this stage produces
__valid_from = 2026-09-25 11:10:00, __valid_to = null, and
__is_current = true. Closing existing history is the destination writer's
responsibility. Other load types skip this transformer. Keep
scd2_effective_column configured for SCD2; without it this stage adds no SCD2
columns.
SystemColumnAdder¶
This class automatically adds audit columns at order 70 on every driver-managed dataflow. There is no metadata field to enable it or define these columns:
| Column | Content |
|---|---|
__created_at |
Framework timestamp when the row was first written |
__updated_at |
Framework timestamp of the current write |
__updated_by |
Configured audit author |
__dataflow_run_id |
ID of the ETL execution or replay chunk that produced the current row/version |
You do not configure these. If a merge destination already has __created_at
from a previous run, the engine preserves its original value on the matched row
and sets __updated_at to the current run timestamp.
Retries reuse one __dataflow_run_id; replay chunks have their own chunk IDs.
For SCD2, closing an old version preserves its original run ID while the new
version receives the current ID. System columns cannot be referenced by
additional_columns (order 30), but can be used in partition expressions
(order 80).
For example, a driver-managed input containing only order_id leaves this
stage with order_id, __created_at, __updated_at, __updated_by, and
__dataflow_run_id. The timestamps, author, and run ID come from the framework
execution; no JSON entry is needed in additional_columns. Direct standalone
use of this class only adds __dataflow_run_id when its caller supplies an ID.
PartitionHandler¶
When to use: the destination needs partition columns, optionally derived
from expressions. Put partition_columns in the
Destination block. Each
Partition column item
names a column and may supply an expression:
"destination": {
"partition_columns": [
{"column": "order_date", "expression": "CAST(updated_at AS DATE)"},
{"column": "region"}
]
}
Without an expression, the column must already exist. Expressions run at order 80, so they can use computed and system columns from earlier stages. The destination writer then uses these columns to partition its output. See Destination partitioning for expression portability and further examples.
Existing, derived, and system-column partitions¶
The example above combines a derived date with an existing region column,
which gives a two-level partition layout. An expression can also replace an
existing column with its derived value. To partition by the current ingestion
date, reference the system timestamp added at order 70:
"destination": {
"partition_columns": [
{"column": "etl_date", "expression": "CAST(__updated_at AS DATE)"}
]
}
An empty or omitted partition_columns list skips this class. Destination
storage options, such as date-folder layout or writer-specific partition
behavior, are configured separately; see Destination and load patterns.
DataMasker¶
When to use: structured scalar values need irreversible masking before
they are written. Configure transform.masking_rules; each
Masking rule specifies a
method, target columns, and method-specific options. This class runs at
order 84, after SCD2, system, and partition columns are available, and before
projection and renaming.
The five methods below are items in transform.masking_rules. A column may
appear in only one masking rule. Use separate columns for the alternatives
below, or choose one method for a given column.
redact: replace non-null values with a constant¶
"transform": {
"masking_rules": [
{"method": "redact", "columns": ["email"], "value": "[REDACTED]"},
{"method": "redact", "columns": ["salary"], "value": "0.00"}
]
}
Here email is a string and salary is a decimal column. All non-null values,
including empty strings, become the constant; existing nulls stay null.
The literal must match the column's type at this stage. The same
typed-literal rules used by fill_null
apply, but masking sees the types after schema conversion.
nullify: remove the value¶
Every value becomes null, retaining the column and its datatype. Use
ColumnProjector.drop_columns if the column itself should disappear.
partial: keep a prefix or suffix¶
"transform": {
"masking_rules": [
{"method": "partial", "columns": ["phone"], "keep_end": 4},
{"method": "partial", "columns": ["account_code"], "keep_start": 2, "keep_end": 2, "mask_char": "#"}
]
}
"12345678" becomes "*5678"; "AB123456YZ" becomes "AB#YZ".
The hidden segment becomes exactly one mask character, not a character per
hidden position. Both keep counts default to zero and must be non-negative;
mask_char defaults to "*" and must be one character. Set both counts to
zero to mask a non-empty string completely.
Only string columns are supported. Null and empty string stay unchanged.
A non-empty value no longer than keep_start + keep_end becomes exactly one
mask character, preventing short values from passing through unchanged.
numeric_bucket: reduce numeric precision¶
The result is floor(value / bucket_size) * bucket_size, cast back to the
original numeric type: 37 becomes 30, and -3 becomes -10. The bucket
size must be positive. Use a size appropriate to the column's datatype;
null stays null, and strings must be converted to numeric before this stage.
date_truncate: reduce date or timestamp precision¶
"transform": {
"masking_rules": [
{"method": "date_truncate", "columns": ["birth_date"], "unit": "year"},
{"method": "date_truncate", "columns": ["invoice_date"], "unit": "month"},
{"method": "date_truncate", "columns": ["event_time"], "unit": "day"},
{"method": "date_truncate", "columns": ["received_at"], "unit": "hour"}
]
}
| Unit | Example input | Result |
|---|---|---|
year |
Date 2026-09-25 |
Date 2026-01-01 |
month |
Date 2026-09-25 |
Date 2026-09-01 |
day |
Timestamp 2026-09-25 11:37:42 |
Timestamp 2026-09-25 00:00:00 |
hour |
Timestamp 2026-09-25 11:37:42 |
Timestamp 2026-09-25 11:00:00 |
These are the supported units. The datatype is preserved; hour requires a
timestamp, and day on an existing Date leaves its date unchanged. String dates
need schema conversion first. Null stays null.
Protected columns and missing inputs¶
Masking merge keys, partition columns, and framework-reserved columns is
rejected. This is column-level PII masking, not dataset anonymization. Keep the
default missing_column_policy = "error" in
transform.configure
for PII-sensitive pipelines so schema drift cannot silently bypass a masking
rule. An unkeyed hash_columns value is
not a substitute for masking low-entropy PII.
An empty or omitted masking_rules list performs no masking.
ColumnProjector¶
When to use: the destination needs only part of the shaped data or different
business-column names. The projector runs at order 85, after masking.
Configure select_columns, drop_columns, and rename_columns in
Transform.
select_columns and drop_columns are mutually exclusive. Selection/removal
uses pre-rename names; rename_columns then renames atomically.
Keep and reorder business columns¶
List the business columns in the desired output order:
Include every merge and partition key. Framework trailing columns are retained automatically and placed after the selected business columns.
Drop unwanted columns¶
All other columns remain. Do not configure select_columns and drop_columns
together.
Rename business columns¶
Rename targets must be distinct and must not overwrite other existing columns.
Chains (a to b, then b to c), swaps, and no-op renames are rejected.
Final names still pass through ColumnNameSanitizer afterwards.
Combine selection/removal with renaming¶
You may combine either select_columns or drop_columns with renaming in the
same block. Resolve the list using names before renaming. If DataMasker
already masked phone, you can make that explicit in the output name:
Projection preserves framework trailing columns and rejects removal or
renaming of merge, partition, and framework-reserved columns. Missing
configured columns fail by default; set
transform.configure
missing_column_policy to ignore only when skipping missing inputs is safe.
Omitting all three projection fields leaves the business columns unchanged.
ColumnNameSanitizer¶
After all configured transforms run, ColumnNameSanitizer applies the
driver's column_name_mode: lower by default, or snake when requested.
This is a runtime option, so there is no field for it in metadata:
# Choose one mode for a run.
result = driver.run(stage="bronze2silver", column_name_mode="lower") # default
Or use word-boundary conversion:
| Input name | lower |
snake |
|---|---|---|
CustomerID |
customerid |
customer_id |
HTTPStatus |
httpstatus |
http_status |
order-date |
order_date |
order_date |
123code |
_123code |
_123code |
__created_at |
unchanged | unchanged |
Both modes clean special characters. Names starting with _ are preserved;
if cleanup would produce an empty name, the original name is retained. A name
collision such as order-date and order date fails instead of creating two
order_date columns. Resolve it earlier with rename_columns or drop_columns.
This class also moves known system/file-info columns to the end. There is no
metadata switch to disable it; already-clean names remain unchanged.
Plan for normalized destination columns when the source uses mixed case,
quoted identifiers, or API keys like CustomerID. The normalized output name
must still match every downstream partition and merge field.
Transform configure flags¶
The reference collects these flags under dataflows[].transform.configure.
These four settings in transform.configure change built-in behavior:
| Key | Default | Effect |
|---|---|---|
convert_timestamp_ntz |
false |
Converts timestamp_ntz columns to timestamp after schema hints when timestamp_timezone is provided |
timestamp_timezone |
unset | Timezone used for an explicit NTZ-to-instant conversion |
deduplicate_by_rank |
false |
Uses rank-based deduplication instead of row-number semantics |
missing_column_policy |
error |
Controls absent columns in typed value/hash/masking rules and projection; schema hints warn and skip, while dedup remains strict |
The class sections above show timestamp and rank cases. Missing-column policy applies across several classes:
Missing-column policy¶
Use error (the default) when every configured input must exist:
"transform": {
"value_rules": [{"operation": "trim", "columns": ["name", "legacy_name"]}],
"configure": {"missing_column_policy": "error"}
}
If legacy_name is absent, this fails. For an optional legacy field, use:
"transform": {
"value_rules": [{"operation": "trim", "columns": ["name", "legacy_name"]}],
"configure": {"missing_column_policy": "ignore"}
}
This still trims name and skips the absent legacy_name.
| Class | What ignore does |
|---|---|
ColumnValueTransformer, DataMasker |
Apply a rule to its existing target columns |
HashColumnAdder |
Skip the whole hash definition if any input is missing; never hash a partial key |
ColumnProjector |
Skip absent select/drop/rename source references |
SchemaConverter |
Independent policy: missing hinted columns warn and skip |
Deduplicator |
Independent policy: missing effective key/order columns fail |
ColumnAdder, RowFilter, PartitionHandler |
This flag does not make SQL expressions tolerate missing inputs |
ignore does not bypass datatype validation, protected-column checks, invalid
expressions, or name collisions. Keep error when skipping a mask would allow
sensitive data through. With no configured operation, each optional class
skips its work; framework system-column handling and sanitization still run.
Full example: multi-pattern transform block¶
This example normalizes an email before masking it, casts the order ID before
hashing it, then renames the masked email before writing. It assumes the source
includes order_id, email, amount, order_date, and updated_at.
{
"name": "orders_to_silver",
"stage": "bronze2silver",
"source": {
"connection_name": "bronze",
"schema_name": "sales",
"table": "orders",
"watermark_columns": ["updated_at"]
},
"destination": {
"connection_name": "silver",
"schema_name": "sales",
"table": "orders",
"load_type": "merge_upsert",
"merge_keys": ["order_id"],
"partition_columns": [
{ "column": "order_date" }
]
},
"transform": {
"value_rules": [
{ "operation": "trim", "columns": ["email"] }
],
"schema_hints": [
{ "column_name": "order_id", "data_type": "long" },
{ "column_name": "amount", "data_type": "decimal", "precision": 18, "scale": 2 },
{ "column_name": "order_date", "data_type": "date" },
{ "column_name": "updated_at", "data_type": "timestamp" }
],
"hash_columns": [
{ "target_column": "order_hash", "columns": ["order_id"], "algorithm": "sha256" }
],
"deduplicate_columns": ["order_id"],
"latest_data_columns": ["updated_at"],
"additional_columns": [
{ "column": "order_year", "expression": "EXTRACT(YEAR FROM order_date)" }
],
"masking_rules": [
{ "method": "partial", "columns": ["email"], "keep_start": 1, "keep_end": 3 }
],
"rename_columns": {"email": "masked_email"},
"configure": {
"convert_timestamp_ntz": false
}
}
}
Common mistakes¶
| Symptom | Likely cause | Fix |
|---|---|---|
| Duplicate rows after merge | No dedup key or no ordering columns | Add deduplicate_columns and latest_data_columns, or let ordering fall back to source.watermark_columns |
| Wrong types in destination table | No schema_hints, or use_schema_hint is disabled on the source connection |
Add hints and confirm source.connection.use_schema_hint is true |
year() function fails on Polars |
Polars SQL doesn't support year() |
Use EXTRACT(YEAR FROM col) |
date() function fails on Polars |
Polars SQL doesn't support date(col) |
Use CAST(col AS DATE) |
| Partition column missing | expression references a column that doesn't exist yet |
Cast or add the column via schema_hints or additional_columns first |
| Hash input column missing | It is only created by additional_columns, which runs after hashing |
Hash existing source columns or calculate the value in a later/custom step |
| Masked field still has its original name | Masking changes values, not names | Use rename_columns after masking_rules |
__updated_at not found in additional_columns |
System columns are added later in the pipeline | Do not reference system columns in additional_columns |
| Destination columns are unexpectedly lowercase | ColumnNameSanitizer runs at the end and defaults to lower |
Expect lowercase output or run with column_name_mode="snake" |