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
- Start connectors with transforms
- Insert test data
- Verify transformations in destination
- 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.