Skip to content

AxonOps — AI-Native Control Plane for Open Source Data Platforms

Kafka Connect Transforms

Single Message Transforms (SMTs) modify records as they flow through Connect without custom code.


Single message transform chain applied to a recordSingle message transform chain applied to a recordTransform ChainSourceRecordTransform 1Transform 2Transform 3FinalRecordTransforms execute in orderEach receives output of previous

TransformDescription
InsertFieldAdd field with static or metadata value
ReplaceFieldInclude, exclude, or rename fields
MaskFieldReplace field value with null
ValueToKeyCopy fields from value to key
ExtractFieldExtract single field from struct
FlattenFlatten nested structures
CastCast field to different type
TransformDescription
RegexRouterRoute to topic based on regex
TimestampRouterRoute based on timestamp
SetSchemaMetadataSet schema name and version
TransformDescription
HeaderFromCopy field to header
InsertHeaderAdd static header
DropHeadersRemove headers
TransformDescription
FilterDrop records matching predicate

{
"name": "my-connector",
"config": {
"connector.class": "...",
"transforms": "transform1,transform2",
"transforms.transform1.type": "...",
"transforms.transform1.field": "...",
"transforms.transform2.type": "...",
"transforms.transform2.regex": "..."
}
}

Apply to key with $Key suffix, value with $Value:

{
"transforms": "addTimestamp",
"transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
"transforms.addTimestamp.timestamp.field": "processed_at"
}

{
"transforms": "addTimestamp",
"transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
"transforms.addTimestamp.timestamp.field": "ingested_at"
}
{
"transforms": "route",
"transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "(.*)_raw",
"transforms.route.replacement": "$1_processed"
}
{
"transforms": "dropFields",
"transforms.dropFields.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
"transforms.dropFields.exclude": "password,ssn,credit_card"
}
{
"transforms": "flatten",
"transforms.flatten.type": "org.apache.kafka.connect.transforms.Flatten$Value",
"transforms.flatten.delimiter": "_"
}

Before: {"user": {"name": "Alice", "age": 30}} After: {"user_name": "Alice", "user_age": 30}

{
"transforms": "cast",
"transforms.cast.type": "org.apache.kafka.connect.transforms.Cast$Value",
"transforms.cast.spec": "price:float64,quantity:int32"
}

{
"name": "events-sink",
"config": {
"connector.class": "...",
"transforms": "addTimestamp,dropSensitive,route",
"transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
"transforms.addTimestamp.timestamp.field": "processed_at",
"transforms.dropSensitive.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
"transforms.dropSensitive.exclude": "password,token",
"transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "events_(.*)",
"transforms.route.replacement": "processed_$1"
}
}

Apply transforms conditionally:

{
"transforms": "insertSource",
"transforms.insertSource.type": "org.apache.kafka.connect.transforms.InsertField$Value",
"transforms.insertSource.static.field": "source",
"transforms.insertSource.static.value": "kafka",
"transforms.insertSource.predicate": "isOrder",
"transforms.insertSource.negate": "false",
"predicates": "isOrder",
"predicates.isOrder.type": "org.apache.kafka.connect.transforms.predicates.TopicNameMatches",
"predicates.isOrder.pattern": "orders-.*"
}

LimitationAlternative
Complex logicCustom SMT or Kafka Streams
Stateful transformsKafka Streams
JoinsKafka Streams
AggregationsKafka Streams