CentralMesh.io

Kafka Connect
AdSense Banner (728x90)

7.2 Common Single Message Transforms

Apply built-in transforms for routing, filtering, field changes, timestamps, and keys.

Common Single Message Transforms

Summary

INSERTFIELD - ADD METADATA

The InsertField SMT adds static values or metadata to messages.

Configuration

json
1{
2  "transforms": "addMetadata",
3  "transforms.addMetadata.type": "org.apache.kafka.connect.transforms.InsertField$Value",
4  "transforms.addMetadata.static.field": "source",
5  "transforms.addMetadata.static.value": "kafka-connect",
6  "transforms.addMetadata.timestamp.field": "ingested_at"
7}

Adds static fields and timestamps for tracking.

REPLACEFIELD - RENAME OR REMOVE

Rename fields to match destination schema or remove sensitive data.

Configuration

json
1{
2  "transforms": "renameFields",
3  "transforms.renameFields.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
4  "transforms.renameFields.renames": "old_name:new_name,user_id:id",
5  "transforms.renameFields.exclude": "internal_field,temp_data"
6}

MASKFIELD - PROTECT SENSITIVE DATA

Replace sensitive values with masked versions.

Configuration

json
1{
2  "transforms": "maskPII",
3  "transforms.maskPII.type": "org.apache.kafka.connect.transforms.MaskField$Value",
4  "transforms.maskPII.fields": "ssn,credit_card",
5  "transforms.maskPII.replacement": "***MASKED***"
6}

TIMESTAMPROUTER - TIME-BASED ROUTING

Route messages to time-based topics for better organization.

Configuration

json
1{
2  "transforms": "routeByTime",
3  "transforms.routeByTime.type": "org.apache.kafka.connect.transforms.TimestampRouter",
4  "transforms.routeByTime.topic.format": "logs-${timestamp}",
5  "transforms.routeByTime.timestamp.format": "yyyy-MM-dd"
6}

Topic logs becomes logs-2024-01-15.

REGEXROUTER - PATTERN-BASED ROUTING

Route based on regex pattern matching.

Configuration

json
1{
2  "transforms": "routeByPattern",
3  "transforms.routeByPattern.type": "org.apache.kafka.connect.transforms.RegexRouter",
4  "transforms.routeByPattern.regex": "(.*)-(.*)",
5  "transforms.routeByPattern.replacement": "$2-$1"
6}

HOISTFIELD - WRAP IN STRUCTURE

Wrap the value in a struct with a single field.

json
1{
2  "transforms": "wrapValue",
3  "transforms.wrapValue.type": "org.apache.kafka.connect.transforms.HoistField$Value",
4  "transforms.wrapValue.field": "payload"
5}

Key ideas

InsertField adds static values and timestamps.

ReplaceField renames or removes fields.

MaskField protects sensitive data.

TimestampRouter creates time-based topic names.

RegexRouter provides pattern-based routing.

HoistField wraps values in structures.