CentralMesh.io

Kafka Connect
AdSense Banner (728x90)

7.4 Hands-On: Single Message Transforms

Build and test a practical multi-transform connector pipeline.

Hands-On: Single Message Transforms

Summary

Let's apply SMTs in a real pipeline, transforming data as it flows from source to sink.

SCENARIO

Stream user data from database to Elasticsearch, adding timestamps and masking emails.

SOURCE CONNECTOR WITH SMTS

json
1{
2  "name": "users-source",
3  "config": {
4    "connector.class": "JdbcSourceConnector",
5    "connection.url": "jdbc:mysql://localhost/db",
6    "table.whitelist": "users",
7    "mode": "incrementing",
8    "incrementing.column.name": "id",
9    "topic.prefix": "db-",
10    "transforms": "addTimestamp,maskEmail,setKey",
11    "transforms.addTimestamp.type": "InsertField$Value",
12    "transforms.addTimestamp.timestamp.field": "ingested_at",
13    "transforms.maskEmail.type": "MaskField$Value",
14    "transforms.maskEmail.fields": "email",
15    "transforms.setKey.type": "ValueToKey",
16    "transforms.setKey.fields": "id"
17  }
18}

TRANSFORMATION FLOW

Original database record:

json
1{
2  "id": 123,
3  "name": "Alice",
4  "email": "[email protected]"
5}

After transforms:

json
1{
2  "id": 123,
3  "name": "Alice",
4  "email": "***@***.com",
5  "ingested_at": 1704067200000
6}

Key set to: {"id": 123}

SINK CONNECTOR WITH SMTS

json
1{
2  "name": "users-sink",
3  "config": {
4    "connector.class": "ElasticsearchSinkConnector",
5    "topics": "db-users",
6    "connection.url": "http://localhost:9200",
7    "transforms": "renameId,routeByDate",
8    "transforms.renameId.type": "ReplaceField$Value",
9    "transforms.renameId.renames": "id:user_id",
10    "transforms.routeByDate.type": "TimestampRouter",
11    "transforms.routeByDate.topic.format": "users-${timestamp}",
12    "transforms.routeByDate.timestamp.format": "yyyy-MM"
13  }
14}

TESTING SMTS

  1. Start connectors with transforms
  2. Insert test data
  3. Verify transformations in destination
  4. Check logs for transformation errors

Key ideas

Apply SMTs to source and/or sink connectors.

Chain multiple transforms for complex logic.

Test with sample data before production.

Monitor for transformation errors.

SMTs eliminate custom transformation code.