Transformations
DataFlow Operator supports message transformations that are applied in list order to each message (transformer 0..N on that message). Across messages, the processor may run transforms in parallel when spec.transformWorkers > 1; sink emit order still matches input order. See Architecture — Pipeline concurrency. Transformations use gjson JSONPath for field access.
Transformation Overview
| Transformation | Description | Input | Output |
|---|---|---|---|
| Timestamp | Adds a timestamp field | 1 message | 1 message |
| Flatten | Expands an array into separate messages | 1 message | N messages |
| Filter | Keeps messages matching a condition (truthy / == / != / metadata.* / && || !) |
1 message | 0 or 1 message |
| Mask | Masks sensitive fields | 1 message | 1 message |
| Router | Sends matching messages to alternate sinks | 1 message | 0 or 1 message |
| Select | Keeps only specified fields | 1 message | 1 message |
| Remove | Removes specified fields | 1 message | 1 message |
| SnakeCase | Converts keys to snake_case | 1 message | 1 message |
| CamelCase | Converts keys to CamelCase | 1 message | 1 message |
| DebeziumUnwrap | Unwraps Debezium envelope into row payload | 1 message | 0 or 1 message |
| ReplaceField | Renames fields; optional include/exclude without flattening | 1 message | 1 message |
| HeadersToPayload | Copies Kafka/message headers into JSON payload fields | 1 message | 1 message |
| StructFlatten | Flattens nested JSON objects into a single-level map | 1 message | 1 message |
| ExtractField | Replaces the payload with the value of one field | 1 message | 1 message |
| HoistField | Wraps the entire payload under a single top-level key | 1 message | 1 message |
| Cast | Converts field values to target types (string / int64 / float64 / bool / null) |
1 message | 1 message (or skip on cast error) |
| Timezone | Converts temporal fields to a target IANA timezone or UTC offset | 1 message | 1 message (or skip on parse error) |
| InsertField | Inserts or overwrites fields with literals, ${metadata.*}, ${now}, or json:<raw> |
1 message | 1 message |
Condition syntax (filter, router, when)
Shared DSL for filter.config.condition, router route conditions, and transformations[].when:
- Truthiness:
$.field/field— true if the payload field exists and is truthy (booleantrue, non-empty string, non-zero number). - Metadata:
metadata.<key>— same checks against messageMetadata(e.g. Kafkatopic,partition,offset). Use$.metadata...only for a JSON field literally namedmetadata. - Equality / inequality:
==/!=with quoted ('value'/"value") or unquoted literals (true,5). - Combinators:
&&(and),||(or;&&binds tighter), unary!, and parentheses(...).
Missing fields fail both truthiness and comparisons. Empty condition is false.
Examples: $.status == 'active' && metadata.topic == 'orders', $.payload.op == 'u' || $.payload.op == 'c', !($.level == 'debug').
Conditional transforms (when)
Optional predicate on any transformation step. If when is set and evaluates to false, this step is skipped and the message continues unchanged (passthrough). Unlike filter, a false when does not drop the message — analogous to a Kafka Connect SMT predicate.
transformations:
- type: debeziumUnwrap
when: "$.payload.op == 'u' || $.payload.op == 'c'"
config:
addOperationFields: true
- type: mask
when: "metadata.topic == 'dbserver1.public.orders'"
config:
fields: ["payload.after.email"]
maskChar: "*"
Timestamp
Adds a timestamp field to each message. Useful for tracking message processing time.
Configuration
transformations:
- type: timestamp
config:
# Field name for timestamp (optional, default: created_at)
fieldName: created_at
# Timestamp format (optional, default: RFC3339)
format: RFC3339
Format
The format value is a Go time layout string. Default is RFC3339 (e.g. 2006-01-02T15:04:05Z07:00). Examples: RFC3339, RFC3339Nano, or custom layouts like 2006-01-02 15:04:05.
Examples
Basic usage
transformations:
- type: timestamp
config:
fieldName: processed_at
Input message:
{
"id": 1,
"name": "Test"
}
Output message:
{
"id": 1,
"name": "Test",
"processed_at": "2024-01-15T10:30:00Z"
}
Custom format
transformations:
- type: timestamp
config:
fieldName: timestamp
format: "2006-01-02 15:04:05"
Output message:
{
"id": 1,
"timestamp": "2024-01-15 10:30:00"
}
Unix timestamp
transformations:
- type: timestamp
config:
fieldName: unix_time
format: Unix
Output message:
{
"id": 1,
"unix_time": "1705312200"
}
Flatten
Expands an array into separate messages, preserving all other fields from the original message. Each array element is merged into the root; objects are flattened to top-level keys. If the field is not an array, the message is returned unchanged. Supports Avro-style arrays wrapped in an object with an array key.
Configuration
transformations:
- type: flatten
config:
# JSONPath to the array to expand (required)
field: items
Examples
Simple flatten
transformations:
- type: flatten
config:
field: items
Input message:
{
"order_id": 12345,
"customer": "John Doe",
"items": [
{"product": "Apple", "quantity": 5},
{"product": "Banana", "quantity": 3}
]
}
Output messages:
{
"order_id": 12345,
"customer": "John Doe",
"product": "Apple",
"quantity": 5
}
{
"order_id": 12345,
"customer": "John Doe",
"product": "Banana",
"quantity": 3
}
Nested arrays
transformations:
- type: flatten
config:
field: orders.items
Input message:
{
"customer_id": 100,
"orders": {
"items": [
{"sku": "SKU001", "price": 10.99},
{"sku": "SKU002", "price": 5.99}
]
}
}
Output messages:
{
"customer_id": 100,
"orders": {},
"sku": "SKU001",
"price": 10.99
}
Combined with other transformations
transformations:
- type: flatten
config:
field: rowsStock
- type: timestamp
config:
fieldName: created_at
This creates a separate message for each array element, each with an added timestamp.
Filter
Keeps only messages that match the condition. Others are dropped.
Condition syntax
See Condition syntax (same DSL as router and when).
Configuration
transformations:
- type: filter
config:
# JSONPath truthiness, or comparison with == / != (required)
condition: "$.status == 'active'"
JSONPath
Uses the gjson library: $.field, $.nested.field, $.array[0], etc. The $. prefix is optional.
Examples
Simple filtering (truthiness)
transformations:
- type: filter
config:
condition: "$.active"
Input messages:
{"id": 1, "active": true} // ✅ Passes
{"id": 2, "active": false} // ❌ Filtered out
{"id": 3} // ❌ Filtered out (no field)
Filtering by equality
transformations:
- type: filter
config:
condition: "$.status == 'active'"
Input messages:
{"status": "active"} // ✅ Passes
{"status": "inactive"} // ❌ Filtered out
{"status": ""} // ❌ Filtered out
{} // ❌ Filtered out (missing field)
Filtering by inequality
transformations:
- type: filter
config:
condition: "$.status != 'deleted'"
Input messages:
{"status": "active"} // ✅ Passes
{"status": "deleted"} // ❌ Filtered out
{} // ❌ Filtered out (missing field)
Filtering by nested path
transformations:
- type: filter
config:
condition: "$.user.status == 'active'"
Input message:
{
"user": {
"status": "active"
}
}
Result: Message passes when user.status equals active.
Mask
Masks sensitive data in specified fields. Supports preserving length or full character replacement.
Configuration
transformations:
- type: mask
config:
# List of JSONPath expressions to fields to mask (required)
fields:
- password
- email
# Character for masking (optional, default: *)
maskChar: "*"
# Preserve original length (optional, default: false)
keepLength: true
Examples
Masking with length preservation
transformations:
- type: mask
config:
fields:
- password
- email
keepLength: true
Input message:
{
"id": 1,
"username": "john",
"password": "secret123",
"email": "john@example.com"
}
Output message:
{
"id": 1,
"username": "john",
"password": "*********",
"email": "****************"
}
Masking with fixed length
transformations:
- type: mask
config:
fields:
- password
keepLength: false
maskChar: "X"
Input message:
{
"password": "verylongpassword123"
}
Output message:
{
"password": "XXX"
}
Masking nested fields
transformations:
- type: mask
config:
fields:
- user.password
- payment.cardNumber
keepLength: true
Input message:
{
"user": {
"password": "secret"
},
"payment": {
"cardNumber": "1234567890123456"
}
}
Output message:
{
"user": {
"password": "******"
},
"payment": {
"cardNumber": "****************"
}
}
Router
Routes messages to different sinks based on conditions. The first matching route determines the sink; if none matches, the message goes to the main sink.
Condition syntax
See Condition syntax (same DSL as filter and when).
Conditions are evaluated in order; the first match wins.
Configuration
transformations:
- type: router
config:
routes:
- condition: "$.level == 'error'"
sink:
type: kafka
config:
brokers: ["localhost:9092"]
topic: error-topic
- condition: "$.level == 'warning'"
sink:
type: postgresql
config:
connectionString: "postgres://..."
table: warnings
Features
- Conditions are checked in order; the first match determines the sink
- If no condition matches, the message goes to the main sink
Examples
Routing by log level
transformations:
- type: router
config:
routes:
- condition: "$.level"
sink:
type: kafka
config:
brokers: ["localhost:9092"]
topic: error-logs
Input messages:
{"level": "error", "message": "Critical error"} // → error-logs topic
{"level": "info", "message": "Info message"} // → main sink
{"level": "warning", "message": "Warning"} // → main sink
Multiple routes
transformations:
- type: router
config:
routes:
- condition: "$.type"
sink:
type: kafka
config:
brokers: ["localhost:9092"]
topic: events-topic
- condition: "$.priority"
sink:
type: postgresql
config:
connectionString: "postgres://..."
table: high_priority_events
Input messages:
{"type": "event", "data": "..."} // → events-topic
{"priority": "high", "data": "..."} // → high_priority_events table
{"data": "..."} // → main sink
Combined with other transformations
transformations:
- type: timestamp
config:
fieldName: processed_at
- type: router
config:
routes:
- condition: "$.level"
sink:
type: kafka
config:
brokers: ["localhost:9092"]
topic: errors
Timestamp is added first, then the message is routed.
Select
Keeps only the specified fields; all others are dropped. Each field is taken by JSONPath; the last path segment is used as the key in the output (e.g. user.name → key name), so the result is flat.
Configuration
transformations:
- type: select
config:
# List of JSONPath expressions for fields to keep (required)
fields:
- id
- name
- email
Examples
Simple field selection
transformations:
- type: select
config:
fields:
- id
- name
- email
Input message:
{
"id": 1,
"name": "John Doe",
"email": "john@example.com",
"password": "secret",
"internal_id": 999
}
Output message:
{
"id": 1,
"name": "John Doe",
"email": "john@example.com"
}
Selecting nested fields (result is flat)
transformations:
- type: select
config:
fields:
- user.id
- user.name
- metadata.timestamp
Input message:
{
"user": {
"id": 1,
"name": "John",
"email": "john@example.com"
},
"metadata": {
"timestamp": "2024-01-15T10:30:00Z",
"source": "api"
}
}
Output message (keys are last path segments):
{
"id": 1,
"name": "John",
"timestamp": "2024-01-15T10:30:00Z"
}
Remove
Removes specified fields from a message. Useful for data cleanup before sending.
Configuration
transformations:
- type: remove
config:
# List of JSONPath expressions to fields to remove (required)
fields:
- password
- internal_id
Examples
Removing sensitive fields
transformations:
- type: remove
config:
fields:
- password
- creditCard
- ssn
Input message:
{
"id": 1,
"name": "John Doe",
"password": "secret",
"creditCard": "1234-5678-9012-3456",
"ssn": "123-45-6789"
}
Output message:
{
"id": 1,
"name": "John Doe"
}
Removing nested fields
transformations:
- type: remove
config:
fields:
- user.password
- metadata.internal
Input message:
{
"user": {
"id": 1,
"name": "John",
"password": "secret"
},
"metadata": {
"timestamp": "2024-01-15",
"internal": "secret"
}
}
Output message:
{
"user": {
"id": 1,
"name": "John"
},
"metadata": {
"timestamp": "2024-01-15"
}
}
Order of Application
Transformations are applied sequentially in the order specified in the transformations list. Each transformation receives the result of the previous one.
Example sequence
transformations:
# 1. Expand array
- type: flatten
config:
field: items
# 2. Add timestamp
- type: timestamp
config:
fieldName: created_at
# 3. Filter inactive
- type: filter
config:
condition: "$.active"
# 4. Remove internal fields
- type: remove
config:
fields:
- internal_id
- debug_info
# 5. Select only needed fields
- type: select
config:
fields:
- id
- name
- created_at
Recommended Order
- Flatten should be first if you need to expand arrays
- Filter apply early to reduce the volume of processed data
- SnakeCase/CamelCase apply after Select/Remove, but before sending to sink
- Mask/Remove apply before Select for security
- Select apply at the end for final cleanup
- Timestamp can be applied anywhere, but usually at the beginning or end
- Router usually applied at the end, after all other transformations
Combined examples
Order processing
transformations:
# Expand items into separate messages
- type: flatten
config:
field: items
# Add timestamp
- type: timestamp
config:
fieldName: processed_at
# Filter paid orders only
- type: filter
config:
condition: "$.status == 'paid'"
# Remove sensitive data
- type: remove
config:
fields:
- customer.creditCard
- customer.cvv
Log processing
transformations:
# Add timestamp
- type: timestamp
config:
fieldName: timestamp
# Mask IP addresses
- type: mask
config:
fields:
- ip_address
keepLength: true
# Route errors
- type: router
config:
routes:
- condition: "$.level"
sink:
type: kafka
config:
brokers: ["localhost:9092"]
topic: error-logs
Field name normalization
transformations:
# Select needed fields
- type: select
config:
fields:
- firstName
- lastName
- email
- address
# Convert to snake_case for PostgreSQL
- type: snakeCase
config:
deep: true
Input message:
{
"firstName": "John",
"lastName": "Doe",
"email": "john@example.com",
"address": {
"streetName": "Main St",
"zipCode": "12345"
}
}
Output message:
{
"first_name": "John",
"last_name": "Doe",
"email": "john@example.com",
"address": {
"street_name": "Main St",
"zip_code": "12345"
}
}
JSONPath support
All transformations that work with fields support JSONPath syntax:
$.field— root field$.nested.field— nested field$.array[0]— array element by index$.array[*]— all array elements$.*— all root-level fields
Performance
- Filter — apply early to reduce data volume
- Select — reduces message size and improves performance
- Flatten — can increase message count; use with care
- Router — creates additional connections; minimize the number of routes
SnakeCase
Converts all JSON object keys to snake_case format. Useful for normalizing field names when integrating with systems using snake_case (e.g., PostgreSQL, Python API).
Configuration
transformations:
- type: snakeCase
config:
# Recursively convert nested objects (optional, default: false)
deep: true
Examples
Simple conversion
transformations:
- type: snakeCase
config:
deep: false
Input message:
{
"firstName": "John",
"lastName": "Doe",
"userName": "johndoe",
"isActive": true,
"itemCount": 42
}
Output message:
{
"first_name": "John",
"last_name": "Doe",
"user_name": "johndoe",
"is_active": true,
"item_count": 42
}
Recursive conversion
transformations:
- type: snakeCase
config:
deep: true
Input message:
{
"firstName": "John",
"address": {
"streetName": "Main St",
"houseNumber": 123,
"zipCode": "12345"
},
"items": [
{
"itemName": "Product",
"itemPrice": 99.99
}
]
}
Output message:
{
"first_name": "John",
"address": {
"street_name": "Main St",
"house_number": 123,
"zip_code": "12345"
},
"items": [
{
"item_name": "Product",
"item_price": 99.99
}
]
}
PascalCase conversion
transformations:
- type: snakeCase
config:
deep: false
Input message:
{
"FirstName": "John",
"LastName": "Doe",
"UserID": 123
}
Output message:
{
"first_name": "John",
"last_name": "Doe",
"user_id": 123
}
Features
- Converts
camelCase→snake_case - Converts
PascalCase→snake_case - Handles consecutive capitals (e.g.
XMLHttpRequest→xml_http_request) - Leaves existing snake_case keys unchanged
- With
deep: falseconverts only top-level keys - With
deep: truerecursively converts all nested objects and arrays
CamelCase
Converts all JSON object keys to CamelCase (PascalCase) format. Useful for normalizing field names when integrating with systems using CamelCase (e.g., Java, C# API).
Configuration
transformations:
- type: camelCase
config:
# Recursively convert nested objects (optional, default: false)
deep: true
Examples
Simple conversion
transformations:
- type: camelCase
config:
deep: false
Input message:
{
"first_name": "John",
"last_name": "Doe",
"user_name": "johndoe",
"is_active": true,
"item_count": 42
}
Output message:
{
"FirstName": "John",
"LastName": "Doe",
"UserName": "johndoe",
"IsActive": true,
"ItemCount": 42
}
Recursive conversion
transformations:
- type: camelCase
config:
deep: true
Input message:
{
"first_name": "John",
"address": {
"street_name": "Main St",
"house_number": 123,
"zip_code": "12345"
},
"items": [
{
"item_name": "Product",
"item_price": 99.99
}
]
}
Output message:
{
"FirstName": "John",
"Address": {
"StreetName": "Main St",
"HouseNumber": 123,
"ZipCode": "12345"
},
"Items": [
{
"ItemName": "Product",
"ItemPrice": 99.99
}
]
}
Single word conversion
transformations:
- type: camelCase
config:
deep: false
Input message:
{
"name": "John",
"id": 123
}
Output message:
{
"Name": "John",
"Id": 123
}
Features
- Converts
snake_case→CamelCase - All words start with capital letter (PascalCase)
- Leaves existing CamelCase keys unchanged
- With
deep: falseconverts only top-level keys - With
deep: truerecursively converts all nested objects and arrays
DebeziumUnwrap
Converts Debezium Kafka events (payload.before/after/op) into a flat row-style message that can be processed by regular transformations and sinks.
Rough analogue of Debezium SMT ExtractNewRecordState (NRSE) for JSON envelopes (not schema-aware Connect records).
Configuration
transformations:
- type: debeziumUnwrap
config:
# Convert Kafka tombstones (empty value) into operation=delete using metadata.key JSON (optional, default: false)
inferDeleteFromTombstone: true
# Copy payload.source.* into metadata as source_<field> (optional, default: false)
includeSourceInMetadata: true
# Operation for Debezium snapshot records (op="r"): insert (default) or update
snapshotOperation: insert
# Write NRSE-style __op / __deleted into the row payload (optional, default: false)
addOperationFields: true
# Copy selected payload.source keys into the row as source_<key> (optional, default: off)
addSourceFields: [table, lsn, ts_ms]
Behavior
op=c-> body frompayload.after,metadata.operation=insertop=u-> body frompayload.after,metadata.operation=updateop=r-> body frompayload.after,metadata.operation=insert(orupdateifsnapshotOperation: update)op=d-> body frompayload.before,metadata.operation=delete- If a message does not contain Debezium envelope (
payload.op), it is passed through unchanged. metadata.operationremainsinsert/update/deleteeven whenaddOperationFieldsis enabled (soft-delete sinks keep reading metadata).
When addOperationFields: true:
- row gets
__op= original Debezium op (c/u/d/r) - row gets
__deleted="true"for deletes (including tombstone-inferred), otherwise"false"(string form, JDBC-friendly)
When addSourceFields is non-empty:
- for each listed key
k, ifpayload.source[k]exists, writesource_<k>into the row - missing keys are skipped; this is independent of
includeSourceInMetadata
Parity vs Debezium ExtractNewRecordState
| Feature | Debezium NRSE | DataFlow debeziumUnwrap |
|---|---|---|
| Flatten envelope to row | Yes | Yes (after / before by op) |
__op / __deleted in payload |
Yes (adds) | Opt-in via addOperationFields |
| Source fields in payload | add.source.fields |
Opt-in via addSourceFields |
| Source in message headers/metadata | Headers | Opt-in via includeSourceInMetadata |
| Drop tombstones | drop.tombstones |
Default drop; keep via inferDeleteFromTombstone (+ optional filter) |
| MongoDB unwrap | Separate SMT | Not supported |
| Schema-aware Connect records | Yes | JSON envelope only |
| Soft-delete column rewrite | Config flags | Use sink softDeleteColumn + metadata.operation |
Kafka tombstones
When inferDeleteFromTombstone: true and message value is empty, the transformer tries to parse metadata.key as JSON (typically Debezium key) and produce a delete message.
If key parsing fails or key is missing, the message is dropped (without error).
With addOperationFields: true, inferred tombstones get __op: "d" and __deleted: "true".
PostgreSQL sink compatibility
For operation=delete, the current PostgreSQL sink uses soft delete via softDeleteColumn. Enable softDeleteColumn in sink.config for delete handling.
Debezium -> PostgreSQL example
transformations:
- type: debeziumUnwrap
config:
inferDeleteFromTombstone: true
includeSourceInMetadata: true
snapshotOperation: insert
addOperationFields: true
addSourceFields: [table, lsn]
ReplaceField
Renames fields and optionally filters the payload via include / exclude. Unlike select, include preserves nested structure (no key flattening). Compatible with Kafka Connect ReplaceField migration patterns.
Configuration
transformations:
- type: replaceField
config:
# Rename mappings oldPath:newPath (optional)
renames:
- oldName:newName
- key.sku:sku
# Keep only these paths, preserving nesting (optional; mutually exclusive with exclude)
include:
- user.id
- status
# Remove these paths (optional; mutually exclusive with include)
# exclude:
# - password
At least one of renames, include, or exclude is required. include and exclude cannot be set together.
Examples
Rename fields
transformations:
- type: replaceField
config:
renames:
- oldName:newName
- key.sku:sku
Input message:
{
"oldName": "value",
"key": { "sku": "12345" },
"other": "keep"
}
Output message:
{
"newName": "value",
"sku": "12345",
"other": "keep"
}
Include without flattening
transformations:
- type: replaceField
config:
include:
- user.id
- user.name
- status
Input message:
{
"user": { "id": 1, "name": "John", "email": "john@example.com" },
"status": "active",
"extra": "drop-me"
}
Output message:
{
"user": { "id": 1, "name": "John" },
"status": "active"
}
HeadersToPayload
Copies Kafka (or other source) message headers from Metadata["headers"] into JSON payload fields. Useful for propagating tracing IDs, tenant keys, and other header values into the body before sinks that do not preserve Kafka headers.
The Kafka source populates Metadata["headers"] as a map[string]string when records have headers.
Configuration
transformations:
- type: headersToPayload
config:
# headerName:fieldPath mappings (required; at least one)
mappings:
- X-Request-Id:requestId
- X-Language:metadata.language
Missing headers are skipped (payload left unchanged for that mapping). Field paths support nested keys and optional $. prefix.
Examples
Copy headers into payload
transformations:
- type: headersToPayload
config:
mappings:
- X-Request-Id:requestId
- X-Language:metadata.language
Input message:
{
"data": "value"
}
Headers: X-Request-Id=req-123, X-Language=en
Output message:
{
"data": "value",
"requestId": "req-123",
"metadata": {
"language": "en"
}
}
StructFlatten
Flattens nested JSON objects into a single-level map (1→1). Unlike flatten (which expands an array into N messages), structFlatten keeps one message and joins nested key paths with a delimiter. Arrays are preserved as values and are not indexed. Compatible with Kafka Connect Flatten migration patterns.
Non-object roots (arrays, primitives) and non-JSON payloads pass through unchanged. Empty nested objects {} produce no keys. Nesting deeper than 64 levels returns a transform error.
Configuration
transformations:
- type: structFlatten
config:
# Delimiter between nested key segments (optional, default: ".")
delimiter: "."
An empty config: {} is valid and uses ".". An explicitly empty delimiter: "" is rejected.
Examples
Dot delimiter (default)
transformations:
- type: structFlatten
config:
delimiter: "."
Input message:
{
"content": {
"id": 42,
"name": {
"first": "David",
"middle": null,
"last": "Wong"
},
"tags": ["a", "b"]
},
"active": true
}
Output message:
{
"content.id": 42,
"content.name.first": "David",
"content.name.middle": null,
"content.name.last": "Wong",
"content.tags": ["a", "b"],
"active": true
}
Underscore delimiter (JDBC / Avro-friendly names)
transformations:
- type: structFlatten
config:
delimiter: "_"
Output keys: content_id, content_name_first, content_name_middle, content_name_last, content_tags, active.
Typical CDC chain: debeziumUnwrap → structFlatten → snakeCase.
ExtractField
Replaces the message payload with the value of a single field (Kafka Connect ExtractField$Value style). Cardinality is always 1→1. Metadata is preserved. Non-JSON payloads and missing paths are passed through unchanged. The new root may be an object, array, primitive, or JSON null.
Configuration
transformations:
- type: extractField
config:
# JSONPath to the field that becomes the new root (required)
field: payload.after # or $.payload.after
Examples
Unwrap nested payload
transformations:
- type: extractField
config:
field: payload.after
Input message:
{"payload":{"after":{"id":1}}}
Output message:
{"id":1}
Extract a primitive or array
transformations:
- type: extractField
config:
field: items
{"items":[1,2,3]} → [1,2,3]
Typical chain before flattening/casting: extractField → structFlatten → cast.
HoistField
Wraps the entire JSON payload under a single top-level key (inverse of extractField). Useful for normalizing envelopes before/after CDC transforms. Non-JSON payloads are unchanged. The wrapper key must be a simple name without dots (not a JSONPath).
Configuration
transformations:
- type: hoistField
config:
# Top-level wrapper key (required; no dots)
field: record
Examples
Wrap a row
transformations:
- type: hoistField
config:
field: record
Input message:
{"id":1}
Output message:
{"record":{"id":1}}
Round-trip with extractField
transformations:
- type: hoistField
config:
field: record
- type: extractField
config:
field: record
Restores the original payload.
Cast
Converts field values to declared target types by JSONPath. Useful after debeziumUnwrap / structFlatten before JDBC or other schema-sensitive sinks. Cardinality is 1→1. Metadata is unchanged. Non-JSON payloads are passed through.
Missing paths are skipped (not an error). Failed conversion of an existing value returns a transform error: the processor logs it, increments dataflow_transformer_errors_total, and skips the message (it is not written to the sink). The operator does not restart. On Kafka, the failed offset is not acked by that message; a later successful write may advance the commit past it (skip-on-error, not infinite retry on the same offset).
Configuration
transformations:
- type: cast
config:
# Map of JSONPath → target type (required, non-empty)
# Types: string | int64 | float64 | bool | null
spec:
id: int64
amount: float64
active: bool
note: string
deleted_at: null # force JSON null
Conversion rules
| Target | Accepted inputs | Errors on |
|---|---|---|
string |
scalars (numbers, bools, strings) | object, array, JSON null |
int64 |
integers, whole floats, numeric strings | fractional numbers, non-numeric values |
float64 |
numbers, numeric strings | non-numeric values |
bool |
bool; strings true/false (case-insensitive); numbers 0/1 |
other values |
null |
any existing value → JSON null |
— |
Examples
After unwrap / flatten
transformations:
- type: debeziumUnwrap
- type: cast
config:
spec:
id: int64
amount: float64
active: bool
Input message:
{"id":"1","amount":"9.99","active":"true"}
Output message:
{"id":1,"amount":9.99,"active":true}
Typical chain: debeziumUnwrap → timezone → cast, or extractField → structFlatten → cast.
Timezone
Converts listed temporal fields to a target IANA timezone or fixed UTC offset (±HH:MM). Useful after debeziumUnwrap before JDBC sinks. Cardinality is 1→1. Metadata is unchanged. Non-JSON payloads are passed through.
Does not interact with spec.maintenance.timezone and is distinct from the timestamp transform (which inserts wall-clock time).
Missing fields and JSON null are skipped. Unparseable values return a transform error: the processor logs it, increments dataflow_transformer_errors_total, and skips the message (same skip-on-error semantics as cast).
Configuration
transformations:
- type: timezone
config:
# Target IANA TZ or UTC offset (required)
timezone: Europe/Moscow
# Fields to convert (required, non-empty JSONPaths)
fields: [created_at, updated_at]
# Assumed source TZ when value has no offset (optional, default: UTC)
sourceTimezone: UTC
# Output layout (optional, default: RFC3339Nano). Also: RFC3339, UnixMilli
format: RFC3339
Input forms
| Input | Behavior |
|---|---|
| RFC3339 / RFC3339Nano string (with or without offset) | Parsed; offsetless values use sourceTimezone |
| Epoch number or numeric string | Milliseconds if \|n\| >= 1e12, otherwise seconds |
Missing path / JSON null |
Skipped |
| Other values | Transform error |
Examples
After debeziumUnwrap
transformations:
- type: debeziumUnwrap
- type: timezone
config:
timezone: Europe/Moscow
fields: [created_at, updated_at]
format: RFC3339
- type: cast
config:
spec:
id: int64
Input message:
{"id":1,"created_at":"2024-01-15T12:00:00Z"}
Output after timezone (Europe/Moscow):
{"id":1,"created_at":"2024-01-15T15:00:00+03:00"}
InsertField
Inserts or overwrites JSON fields with static literals, Kafka/record metadata placeholders, wall-clock time, or raw JSON values. Analogous to Kafka Connect InsertField$Value, but with DataFlow metadata and json: literals. Cardinality is 1→1. Metadata is unchanged. Non-JSON payloads are passed through.
Distinct from timestamp (single configurable now field) and headersToPayload (headers → body only).
Configuration
transformations:
- type: insertField
config:
fields:
pipeline: "orders-cdc" # literal string
source_topic: "${metadata.topic}" # metadata placeholder
source_partition: "${metadata.partition}"
source_offset: "${metadata.offset}"
source_timestamp: "${metadata.timestamp}"
ingested_at: "${now}" # wall-clock RFC3339
"flags.reprocessed": "json:false" # raw JSON value
Value syntax
| Value | Result |
|---|---|
| any other string | written as a JSON string literal |
${metadata.<key>} |
string form of Metadata[key]; missing/nil → "" |
${now} |
time.Now() formatted as RFC3339 |
json:<raw> |
parsed JSON written as a native value (false, 42, objects, arrays) |
Common metadata keys from Kafka source: topic, partition, offset, timestamp. Invalid json: content is a transform error (message skipped).
Examples
Enrich with pipeline context
transformations:
- type: insertField
config:
fields:
pipeline: "orders-cdc"
source_topic: "${metadata.topic}"
ingested_at: "${now}"
"flags.reprocessed": "json:false"
- type: snakeCase
config:
deep: true
Input: {"id":1,"Name":"A"} with metadata topic=raw.events
After insertField: fields pipeline, source_topic, ingested_at, flags.reprocessed added
After snakeCase: keys converted to snake_case
Limitations
- Filter / Router /
when: shared condition DSL — truthiness,==/!=,metadata.*,&&/||/!, parentheses. Missing fields fail the condition. Groovy/JS/CEL scripting is not supported.whenskips a step (passthrough);filterdrops the message. - Flatten: works only with arrays (including Avro-style wrapper with
arraykey), not arbitrary objects. For nested-object flattening usestructFlatten. - StructFlatten: 1→1 object flatten; arrays are kept as values (use array
flattenfirst to explode). Max nesting depth is 64. - ExtractField: missing path and non-JSON → passthrough; root may become a non-object (later object-only transforms will passthrough).
- HoistField: wrapper key must be a simple top-level name (no dots); wraps any JSON value including arrays and primitives.
- Cast: missing path → skip; failed conversion of an existing value → transform error (message skipped, not written to sink; Kafka offset may advance past it on later success). Types:
string,int64,float64,bool,null. - Timezone: only listed
fields; missing/null → skip; unparseable → transform error (same skip-on-error as cast). Formats:RFC3339,RFC3339Nano(default),UnixMilli. Not related tomaintenance.timezoneor thetimestamptransform. - InsertField: non-JSON → passthrough; missing metadata → empty string; invalid
json:→ transform error. Distinct fromtimestampandheadersToPayload. - Select: result is always flat; the key is the last JSONPath segment.
- ReplaceField:
includepreserves nesting (unlikeselect);includeandexcludeare mutually exclusive. - HeadersToPayload: requires headers in
Metadata["headers"](Kafka source sets this); missing headers are skipped; non-JSON payloads are unchanged. - SnakeCase and CamelCase: work only with valid JSON; binary data is returned unchanged.
- DebeziumUnwrap: supports Debezium JSON envelope (
payload.op/before/after). Foroperation=deletein PostgreSQL sink,softDeleteColumnis required; otherwise delete events can be treated as regular insert/update writes.