본문 바로가기

server/kafka

Kafka connect (2) - SMT(single message Transform)

지난 글에서 FileStreamSinkConnector로 토픽 메시지를 파일로 저장하는 작업을 진행했다.  이번에는 커넥터에 SMT(Single Message Transform)를 추가해서 메시지를 변형해보는것을 진행한다.  컨슈머 코드 없이 json 설정만으로 필드를 추가, 삭제, 변경한다. 

사용할 트랜스폼은 두가지

 

  • InsertField$Value  :  타임스탬프, 고정값, kafka 메타(topic/partition/offset) 필드 추가
  • ReplaceField$Value : 필드 제거, 필드 이름 변경

 

1. SMT

sink 커넥터 기준으로 메시지 흐름은 다음과 같다. 

kafka topic → converter (bytes → map/struct) → transform 체인 → sink task

 

converter가 바이트를 구조화 된 데이터를 확인하고 SMT는 레코드를 한건씩 수정한다. 

 

참고로 레코드 하나에서 끝나느 가벼운 변형(필드 추가/제거/수정/마스킹) 이 SMT의 역활이며 조인이나 집계처럼 여러 레코드를 봐야 하는 처리는 streams나 ksqldb로 넘겨야 한다. 

 

broker 컨테이너에 진입한 상태에서 실습용 토픽을 만든다.

$ /bin/kafka-topics --bootstrap-server localhost:9092 \
--create --replication-factor 1 --partitions 1 --topic topic-smt


Created topic topic-smt.

 

 

2. InsertField$Value — 타임스탬프, 고정값, kafka 메타 정보 추가하기

 

커넥터 설정을 파일로 작성.

지난 글의 file sink에 transforms 항목을 추가한 형태로 먼저 addTimestamp 부터 시작.

transforms에 적는 이름(addTimestamp)은 내가 정하는 별칭

(아래는 실패로 떨어질 예정)

$ vim /tmp/smt-sink.json 

{
  "connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
  "tasks.max": "1",
  "topics": "topic-smt",
  "file": "/tmp/smt-sink.out",
  "value.converter": "org.apache.kafka.connect.storage.StringConverter",
  "transforms": "addTimestamp",
  "transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addTimestamp.timestamp.field": "event_ts"
}

 

smt sink 등록

# 등록 

$ curl -X PUT -H "Content-Type: application/json" \
-d @/tmp/smt-sink.json \
localhost:8083/connectors/smt-sink/config


# 등록 확인

$ curl localhost:8083/connectors/smt-sink | jq
  % Total    % Received % Xferd  Average Speed   Time    Time     Time  Current
                                 Dload  Upload   Total   Spent    Left  Speed
100   480  100   480    0     0  96000      0 --:--:-- --:--:-- --:--:-- 96000
{
  "name": "smt-sink",
  "config": {
    "connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
    "file": "tmp/smt-sink.out",
    "transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
    "tasks.max": "1",
    "topics": "topic-smt",
    "transforms": "addTimestamp",
    "name": "smt-sink",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "transforms.addTimestamp.timestamp.field": "event_ts"
  },
  "tasks": [
    {
      "connector": "smt-sink",
      "task": 0
    }
  ],
  "type": "sink"
}

 

 

프로듀서로 json 메시지를 한 건 추가 테스트

$ /bin/kafka-console-producer --bootstrap-server localhost:9092 --topic topic-smt
>{"user_id": "u-1", "action": "click", "password": "1234"}

 

상태를 확인해보면 태스크가 죽어 있다.

$ curl localhost:8083/connectors/smt-sink/status | jq
  % Total    % Received % Xferd  Average Speed   Time    Time     Time  Current
                                 Dload  Upload   Total   Spent    Left  Speed
100  2589  100  2589    0     0  78454      0 --:--:-- --:--:-- --:--:-- 78454
{
  "name": "smt-sink",
  "connector": {
    "state": "RUNNING",
    "worker_id": "172.19.0.3:8083"
  },
  "tasks": [
    {
      "id": 0,
      "state": "FAILED",
      "worker_id": "172.19.0.3:8083",
      "trace": "org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler\n\tat org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:230)\n\tat org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:156)\n\tat org.apache.kafka.connect.runtime.TransformationChain.apply(TransformationChain.java:53)\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.convertAndTransformRecord(WorkerSinkTask.java:552)\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:505)\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:341)\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:242)\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:211)\n\tat org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:204)\n\tat org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:259)\n\tat org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:181)\n\tat java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)\n\tat java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)\n\tat java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)\n\tat java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)\n\tat java.base/java.lang.Thread.run(Thread.java:829)\nCaused by: org.apache.kafka.connect.errors.DataException: Only Struct objects supported for [field insertion], found: java.lang.String\n\tat org.apache.kafka.connect.transforms.util.Requirements.requireStruct(Requirements.java:52)\n\tat org.apache.kafka.connect.transforms.InsertField.applyWithSchema(InsertField.java:164)\n\tat org.apache.kafka.connect.transforms.InsertField.apply(InsertField.java:135)\n\tat org.apache.kafka.connect.runtime.TransformationStage.apply(TransformationStage.java:57)\n\tat org.apache.kafka.connect.runtime.TransformationChain.lambda$apply$0(TransformationChain.java:53)\n\tat org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:180)\n\tat org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:214)\n\t... 15 more\n"
    }
  ],
  "type": "sink"
}

 

 

InsertField는 map(또는 schema 있는 struct)에 필드를 넣는 트랜스폼인데, StringConverter는 값을 통짜 문자열로 넘기므로 필드 변경이 불가능하다.  

지난 글에서는 평문을 그대로 파일에 쓰기만 해서 StringConverter로 충분했지만, SMT로 필드를 만지려면 converter가 json을 map으로 풀어줘야 한다.

value.converter를 JsonConverter로 바꾸고 schemas.enable을 꺼야 한다. 

 

$ vim /tmp/smt-sink.json 

{
  "connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
  "tasks.max": "1",
  "topics": "topic-smt",
  "file": "/tmp/smt-sink.out",
  "value.converter": "org.apache.kafka.connect.json.JsonConverter",
  "value.converter.schemas.enable": "false",
  "transforms": "addTimestamp",
  "transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addTimestamp.timestamp.field": "event_ts"
}

 

# 다시 메시지 추가

$ /bin/kafka-console-producer --bootstrap-server localhost:9092 --topic topic-smt
>{"user_id": "u-1", "action": "click", "password": "1234"}
>

 

 

최종 파일을 살펴보면 envet_ts필드가 추가된것을 확인할수 있다. 

$ tail -f /tmp/smt-sink.out
{action=click, password=1234, user_id=u-1, event_ts=1788055347196}

 

이번엔 고정값을 추가해보자 

env : dev로 추가하는거 테스트

$ vim /tmp/smt-sink.json 

{
  "connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
  "tasks.max": "1",
  "topics": "topic-smt",
  "file": "/tmp/smt-sink.out",
  "value.converter": "org.apache.kafka.connect.json.JsonConverter",
  "value.converter.schemas.enable": "false",
  "transforms": "addTimestamp,addEnv",
  "transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addTimestamp.timestamp.field": "event_ts",
  "transforms.addEnv.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addEnv.static.field": "env",
  "transforms.addEnv.static.value": "dev"
}

 

# sink 재 등록 해주고 

$ curl -X PUT -H "Content-Type: application/json" \
-d @/tmp/smt-sink.json \
localhost:8083/connectors/smt-sink/config


# 토픽에 메시지 추가 (1~2초 대기)

$ /bin/kafka-console-producer --bootstrap-server localhost:9092 --topic topic-smt
>{"user_id": "u-2", "action": "view", "password": "1234"}


# 데이터 확인 
# env=dev 입력 완료
$ tail -1 /tmp/smt-sink.out
{action=view, password=1234, env=dev, user_id=u-2, event_ts=1788059117575}

 

 

이번엔 kafka 메타 값 추가

$ vim /tmp/smt-sink.json 

{
"connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
  "tasks.max": "1",
  "topics": "topic-smt",
  "file": "/tmp/smt-sink.out",
  "value.converter": "org.apache.kafka.connect.json.JsonConverter",
  "value.converter.schemas.enable": "false",
  "transforms": "addTimestamp,addEnv,addMeta",
  "transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addTimestamp.timestamp.field": "event_ts",
  "transforms.addEnv.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addEnv.static.field": "env",
  "transforms.addEnv.static.value": "dev",
  "transforms.addMeta.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addMeta.topic.field": "kafka_topic",
  "transforms.addMeta.partition.field": "kafka_partition",
  "transforms.addMeta.offset.field": "kafka_offset"
}

 

 

# sink 적용 
$ curl -X PUT -H "Content-Type: application/json"   -d @/tmp/smt-sink.json   localhost:8083/connectors/smt-sink/config
{"name":"smt-sink","config":{"connector.class":"org.apache.kafka.connect.file.FileStreamSinkConnector","tasks.max":"1","topics":"topic-smt","file":"/tmp/smt-sink.out","value.converter":"org.apache.kafka.connect.json.JsonConverter","value.converter.schemas.enable":"false","transforms":"addTimestamp,addEnv,addMeta","transforms.addTimestamp.type":"org.apache.kafka.connect.transforms.InsertField$Value","transforms.addTimestamp.timestamp.field":"event_ts","transforms.addEnv.type":"org.apache.kafka.connect.transforms.InsertField$Value","transforms.addEnv.static.field":"env","transforms.addEnv.static.value":"dev","transforms.addMeta.type":"org.apache.kafka.connect.transforms.InsertField$Value","transforms.addMeta.topic.field":"kafka_topic","transforms.addMeta.partition.field":"kafka_partition","transforms.addMeta.offset.field":"kafka_offset","name":"smt-sink"},"tasks":[{"connector":"smt-sink","task":0}],"type":"sink"}


# 메시지 추가 
$  /bin/kafka-console-producer --bootstrap-server localhost:9092 --topic topic-smt
>{"user_id": "u-3", "action": "click", "password": "1234"}

# 데이터 확인 

$ tail -1 /tmp/smt-sink.out
{password=1234, kafka_offset=6, kafka_partition=0, user_id=u-3, event_ts=1788059286859, action=click, kafka_topic=topic-smt, env=dev}

 

 

3. ReplaceField$Value — 필드 제거, 이름 변경

password처럼 암호화 없이 메시지를 보내면 안 되는 필드를 제거해본다.

$ vim /tmp/smt-sink.json
{
"connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
  "tasks.max": "1",
  "topics": "topic-smt",
  "file": "/tmp/smt-sink.out",
  "value.converter": "org.apache.kafka.connect.json.JsonConverter",
  "value.converter.schemas.enable": "false",
  "transforms": "addTimestamp,addEnv,addMeta,dropSecret",
  "transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addTimestamp.timestamp.field": "event_ts",
  "transforms.addEnv.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addEnv.static.field": "env",
  "transforms.addEnv.static.value": "dev",
  "transforms.addMeta.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addMeta.topic.field": "kafka_topic",
  "transforms.addMeta.partition.field": "kafka_partition",
  "transforms.addMeta.offset.field": "kafka_offset",
  "transforms.dropSecret.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
  "transforms.dropSecret.exclude": "password",
  "transforms.rename.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
  "transforms.rename.renames": "user_id:uid,action:event_name"
}

 

# sink 적용 
$ curl -X PUT -H "Content-Type: application/json"   -d @/tmp/smt-sink.json   localhost:8083/connectors/smt-sink/config
{"name":"smt-sink","config":{"connector.class":"org.apache.kafka.connect.file.FileStreamSinkConnector","tasks.max":"1","topics":"topic-smt","file":"/tmp/smt-sink.out","value.converter":"org.apache.kafka.connect.json.JsonConverter","value.converter.schemas.enable":"false","transforms":"addTimestamp,addEnv,addMeta","transforms.addTimestamp.type":"org.apache.kafka.connect.transforms.InsertField$Value","transforms.addTimestamp.timestamp.field":"event_ts","transforms.addEnv.type":"org.apache.kafka.connect.transforms.InsertField$Value","transforms.addEnv.static.field":"env","transforms.addEnv.static.value":"dev","transforms.addMeta.type":"org.apache.kafka.connect.transforms.InsertField$Value","transforms.addMeta.topic.field":"kafka_topic","transforms.addMeta.partition.field":"kafka_partition","transforms.addMeta.offset.field":"kafka_offset","name":"smt-sink"},"tasks":[{"connector":"smt-sink","task":0}],"type":"sink"}


# 메시지 추가 
$  /bin/kafka-console-producer --bootstrap-server localhost:9092 --topic topic-smt
>{"user_id": "u-5", "action": "click", "password": "1234"}

# 데이터 확인 
# password 가 삭제 된것을 볼수 잇음


$ tail -1 /tmp/smt-sink.out
{kafka_offset=8, kafka_partition=0, user_id=u-5, event_ts=1788059746302, action=click, kafka_topic=topic-smt, env=dev}

 

 

 

번외. uid 필드 단방향 해시 — aiven Hash SMT

특정 메시지의 필드를 이름만 바꿔서 내리는 걸로는 부족하고 uid를 암호화해서 적재하고 싶을 때가 있다.

kafka 기본 SMT에는 해시가 없고, 커뮤니티 SMT인 aiven transforms-for-apache-kafka-connect의 Hash를 사용해본다.( md5/sha1/sha256을 지원)

 

 

먼저 해당 플러그인 설치 

$ yum install -y unzip

$ cd /tmp
$ curl -sfLO https://github.com/Aiven-Open/transforms-for-apache-kafka-connect/releases/download/v1.5.0/transforms-for-apache-kafka-connect-1.5.0.zip
$ mkdir -p /tmp/plugins
$ unzip -q transforms-for-apache-kafka-connect-1.5.0.zip -d /tmp/plugins


# plugin path 추가 
$ echo "plugin.path=/tmp/plugins" >> /tmp/connect-distributed.cp-kafka-7.5.4.properties

# 워커 재실행 
$ CLASSPATH=/tmp/connect-file-3.5.1.jar /bin/connect-distributed /tmp/connect-distributed.cp-kafka-7.5.4.properties


# 적용 확인 
$ curl "localhost:8083/connector-plugins?connectorsOnly=false" | jq '.[] | select(.type=="transformation") | .class' | grep -i hash


"io.aiven.kafka.connect.transforms.Hash$Key"
"io.aiven.kafka.connect.transforms.Hash$Value"

 

 

다시 sink 수정 

uid로 들어온 필드의 값을 해시값으로 변경 셋팅




$ vim /tmp/smt-sink.json
{
"connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
  "tasks.max": "1",
  "topics": "topic-smt",
  "file": "/tmp/smt-sink.out",
  "value.converter": "org.apache.kafka.connect.json.JsonConverter",
  "value.converter.schemas.enable": "false",
  "transforms": "addTimestamp,addEnv,addMeta,dropSecret,rename,hashUid",
  "transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addTimestamp.timestamp.field": "event_ts",
  "transforms.addEnv.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addEnv.static.field": "env",
  "transforms.addEnv.static.value": "dev",
  "transforms.addMeta.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addMeta.topic.field": "kafka_topic",
  "transforms.addMeta.partition.field": "kafka_partition",
  "transforms.addMeta.offset.field": "kafka_offset",
  "transforms.dropSecret.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
  "transforms.dropSecret.exclude": "password",
  "transforms.rename.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",
  "transforms.rename.renames": "user_id:uid,action:event_name",
  "transforms.hashUid.type": "io.aiven.kafka.connect.transforms.Hash$Value",
  "transforms.hashUid.field.name": "uid",
  "transforms.hashUid.function": "sha256"
}

 

 

 

# 설정 재시작

$ curl -X PUT -H "Content-Type: application/json" \
-d @/tmp/smt-sink.json \
localhost:8083/connectors/smt-sink/config

$ curl -X POST localhost:8083/connectors/smt-sink/tasks/0/restart



# 메시지 추가
# 이번에는 스키마가 필요
$ /bin/kafka-console-producer --bootstrap-server localhost:9092 --topic topic-smt
> {"schema":{"type":"struct","fields":[{"field":"user_id","type":"string"},{"field":"action","type":"string"},{"field":"password","type":"string"}]},"payload":{"user_id":"u-7","action":"click","password":"1234"}}



# 데이터 확인

$ tail -1 /tmp/smt-sink.out
Struct{uid=bf9023e0fc14f272cb98cf9d5c5f3e5cc4dc81206799b0b977c681738d7308b3,event_name=click,event_ts=Sun Aug 30 03:36:29 UTC 2026,env=dev,kafka_topic=topic-smt,kafka_partition=0,kafka_offset=13}

 

 

끝.