Blog

Kafka Connect (1)

Kafka Connect

โ€ข
Kafka Connect ๋Š” ์™ธ๋ถ€ ์‹œ์Šคํ…œ์— ์—ฐ๊ฒฐํ•˜๊ธฐ ์œ„ํ•œ ํ”„๋ ˆ์ž„์›Œํฌ๋ฅผ ์ œ๊ณตํ•˜๋Š” KAFKA ์˜ ์˜คํ”ˆ ์†Œ์Šค ๊ตฌ์„ฑ ์š”์†Œ์ด๋‹ค.
โ€ข
Kafka Connect ๋ฅผ ์‚ฌ์šฉํ•˜๊ธฐ ์œ„ํ•ด์„œ๋Š” ํด๋Ÿฌ์Šคํ„ฐ๋ฅผ ๊ณ„ํšํ•˜๊ณ  ํ”„๋กœ๋น„์ €๋‹ํ•˜๊ณ , ์ž‘์—…์„ ์ฒ˜๋ฆฌํ•˜๊ณ  ๋ถ€ํ•˜์— ๋”ฐ๋ผ ์Šค์ผ€์ผ๋ง์„ ๊ณ ๋ คํ•ด์•ผ ํ•œ๋‹ค.

์šฉ์–ด

โ€ข
Connect: Connector ๋ฅผ ๋™์ž‘ํ•˜๊ฒŒ ํ•˜๋Š” ํ”„๋กœ์„ธ์Šค (์„œ๋ฒ„)
โ€ข
Connector: Data Source(DB) ์˜ ๋ฐ์ดํ„ฐ๋ฅผ ์ฒ˜๋ฆฌํ•˜๋Š” ์†Œ์Šค๊ฐ€ ๋“ค์–ด์žˆ๋Š” jar ํŒŒ์ผ
โ€ข
Source Connector: Producer ์—ญํ• ์„ ๋‹ด๋‹น
โ€ข
Sink Connector: Consumer ์—ญํ• ์„ ๋‹ด๋‹น
โ€ข
Standalone ๋ชจ๋“œ: ํ•˜๋‚˜์˜ Connect ๋งŒ ์‚ฌ์šฉํ•˜๋Š” ๋ชจ๋“œ
โ€ข
Distributed ๋ชจ๋“œ: ์—ฌ๋Ÿฌ๊ฐœ์˜ Connect ๋ฅผ ํ•œ๊ฐœ์˜ ํด๋Ÿฌ์Šคํ„ฐ๋กœ ๋ฌถ์–ด์„œ ์‚ฌ์šฉํ•˜๋Š” ๋ชจ๋“œ
โ—ฆ
Connect ์ค‘ ํ•œ๊ฐœ๊ฐ€ ์žฅ์• ๋‚˜๋„ ๋‚˜๋จธ์ง€ Connect ๋“ค์ด ์ด์–ด์„œ ์ฒ˜๋ฆฌํ•  ์ˆ˜ ์žˆ์Œ

์„ค์น˜ ๋ฐฉ๋ฒ•

โ€ข
์ง์ ‘ ์„ค์น˜ํ•˜๋Š” ๋ฐฉ๋ฒ•๋„ ์žˆ๊ฒ ์ง€๋งŒ, ๋กœ์ปฌ์—์„œ ํ…Œ์ŠคํŠธ๋กœ ํ•„์š”ํ•˜๋ฉด docker ๋กœ ๋„์šธ ์ˆ˜ ์žˆ๋‹ค
docker run -d \ --name=kafka-connect \ -e CONNECT_BOOTSTRAP_SERVERS=localhost:29092 \ -e CONNECT_REST_PORT=28083 \ -e CONNECT_GROUP_ID="quickstart-avro" \ -e CONNECT_CONFIG_STORAGE_TOPIC="quickstart-config" \ -e CONNECT_OFFSET_STORAGE_TOPIC="quickstart-offsets" \ -e CONNECT_STATUS_STORAGE_TOPIC="quickstart-status" \ -e CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR=1 \ -e CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR=1 \ -e CONNECT_STATUS_STORAGE_REPLICATION_FACTOR=1 \ -e CONNECT_KEY_CONVERTER="org.apache.kafka.connect.json.JsonConverter" \ -e CONNECT_VALUE_CONVERTER="org.apache.kafka.connect.json.JsonConverter" \ -e CONNECT_INTERNAL_KEY_CONVERTER="org.apache.kafka.connect.json.JsonConverter" \ -e CONNECT_INTERNAL_VALUE_CONVERTER="org.apache.kafka.connect.json.JsonConverter" \ -e CONNECT_REST_ADVERTISED_HOST_NAME="localhost" \ -e CONNECT_LOG4J_ROOT_LOGLEVEL=DEBUG \ -e CONNECT_PLUGIN_PATH=/usr/share/java,/etc/kafka-connect/jars \ -v /tmp/quickstart/file:/tmp/quickstart \ confluentinc/cp-kafka-connect:latest
Shell
๋ณต์‚ฌ

REST API

โ€ข
Kafka Connect ๋Š” REST API ๋ฅผ ์‚ฌ์šฉํ•ด์„œ Connector ๋ฅผ ๋งŒ๋“ค ์ˆ˜ ์žˆ๋Š” ๋ฐฉ๋ฒ•์„ ์ œ๊ณตํ•œ๋‹ค
โ€ข
GET /connectors : ์ปค๋„ฅํ„ฐ ๋ชฉ๋ก ์กฐํšŒ
โ€ข
GET /connectors/{name} : {name} ์ปค๋„ฅํ„ฐ ์ •๋ณด ์กฐํšŒ
โ€ข
GET /connectors/{name}/status : ์ปค๋„ฅํ„ฐ ์ƒํƒœ ์กฐํšŒ
โ€ข
POST /connectors : ์ปค๋„ฅํ„ฐ ์ƒ์„ฑ
โ€ข
DELETE /connectors/{name} : {name} ์ปค๋„ฅํ„ฐ ์‚ญ์ œ
โ€ข
GET /connector-plugins : ์„ค์น˜๋œ ํ”Œ๋Ÿฌ๊ทธ์ธ ๋ชฉ๋ก ์กฐํšŒ

Debezium

โ€ข
๋ฐ์ดํ„ฐ๋ฒ ์ด์Šค์˜ ๋ณ€๊ฒฝ ์‚ฌํ•ญ์„ ์บก์ฒ˜ํ•˜์—ฌ ์• ํ”Œ๋ฆฌ์ผ€์ด์…˜์—์„œ ์‚ฌ์šฉํ•  ์ˆ˜ ์žˆ๋„๋ก ํ•ด์ฃผ๋Š” ๋ถ„์‚ฐ ์„œ๋น„์Šค ์ง‘ํ•ฉ์ด๋‹ค.
โ€ข
CDC ( Change Data Capture ) ์„ ์‚ฌ์šฉํ•˜์—ฌ ๋ฐ์ดํ„ฐ๋ฒ ์ด์Šค์˜ ๋ณ€๊ฒฝ ์‚ฌํ•ญ์„ ์ˆ˜์ง‘ํ•œ๋‹ค
โ€ข
๋ชจ๋“  Low Level ๋ณ€๊ฒฝ์„ Changed Event Stream ์— ๊ธฐ๋กํ•œ๋‹ค.
โ—ฆ
์• ํ”Œ๋ฆฌ์ผ€์ด์…˜์€ ์ด ์ŠคํŠธ๋ฆผ์„ ํ†ตํ•ด ๋ณ€๊ฒฝ ์ด๋ฒคํŠธ๋“ค์„ ์ˆœ์„œ๋Œ€๋กœ ์ฝ๋Š”๋‹ค
โ€ข
Debezium ์˜ ๋ชฉํ‘œ๋Š” ๋‹ค์–‘ํ•œ DBMS ๋“ค์˜ ๋ณ€๊ฒฝ ์‚ฌํ•ญ์„ ์บก์ฒ˜ํ•˜๊ณ  ๋น„์Šทํ•œ ๊ตฌ์กฐ์˜ ๋ณ€๊ฒฝ ์ด๋ฒคํŠธ๋“ค์„ Producing ํ•˜๋Š” Connector ๋ผ์ด๋ธŒ๋Ÿฌ๋ฆฌ๋ฅผ ๊ตฌ์ถ•ํ•˜๋Š”๊ฒƒ์ด๋‹ค.

์ง€์›ํ•˜๋Š” ์ปค๋„ฅํ„ฐ

โ€ข
MySQL
โ€ข
MongoDB
โ€ข
PostgreSQL
โ€ข
Orcal
โ€ข
SQL Server
โ€ข
DB2
โ€ข
Cassandra
โ€ข
Vitess

ํŠน์ง•

โ€ข
Debezium ์€ Log Based CDC ์ด๋‹ค.
โ€ข
๋ชจ๋“  ๋ฐ์ดํ„ฐ ๋ณ€๊ฒฝ์ด ์บก์ฒ˜๋œ๋‹ค
โ€ข
Data Model ๋ณ€๊ฒฝ์ด ํ•„์š” ์—†๋‹ค
โ€ข
๋ณ€๊ฒฝ ๋ฟ๋งŒ ์•„๋‹ˆ๋ผ ์‚ญ์ œ๋„ ์บก์ฒ˜ํ•œ๋‹ค
โ€ข
๋ ˆ์ฝ”๋“œ์˜ ๊ณผ๊ฑฐ ์ƒํƒœ๋„ ์บก์ฒ˜๊ฐ€ ๊ฐ€๋Šฅํ•˜๋‹ค.

๊ธฐ๋Šฅ

โ€ข
Snapshots: ์ปค๋„ฅํ„ฐ๊ฐ€ ์‹œ์ž‘๋  ๋•Œ ๋ฐ์ดํ„ฐ๋ฒ ์ด์Šค์˜ ํ˜„์žฌ ์ƒํƒœ์— ๋Œ€ํ•œ ์ดˆ๊ธฐ ์Šค๋ƒ…์ƒท์„ ์ƒ์„ฑํ•  ์ˆ˜ ์žˆ๋‹ค.
โ€ข
Filters: ํŠน์ • ํ…Œ์ด๋ธ”์ด๋‚˜ ์ปฌ๋Ÿผ์˜ ๋ณ€๊ฒฝ๋งŒ ์บก์ฒ˜ํ•  ์ˆ˜ ์žˆ๋‹ค
โ€ข
Masking: ๋ฏผ๊ฐํ•œ ์ •๋ณด์˜ ๊ฒฝ์šฐ ํŠน์ • ์ปฌ๋Ÿผ์„ ๋งˆ์Šคํ‚น์ฒ˜๋ฆฌ ํ•  ์ˆ˜ ์žˆ๋‹ค
โ€ข
Message Transformation
โ—ฆ
Topic Routing
โ–ช
์ด ๊ธฐ๋Šฅ์„ ์‚ฌ์šฉํ•˜๋ฉด ํ† ํ”ฝ์˜ ์ด๋ฆ„์„ ๋ณ€๊ฒฝํ•˜๊ฑฐ๋‚˜, ์—ฌ๋Ÿฌ๊ฐœ์˜ ํ…Œ์ด๋ธ”์˜ ๋ณ€๊ฒฝ์‚ฌํ•ญ์„ ํ•˜๋‚˜์˜ ํ† ํ”ฝ์œผ๋กœ ์ „๋‹ฌํ•  ์ˆ˜ ์žˆ๋‹ค
โ—ฆ
New Record State Extraction
โ—ฆ
Outbox Event Router
โ—ฆ
Message Filtering
โ—ฆ
Content-Based Routing

์•„ํ‚คํ…์ฒ˜

Kafka Connect

โ€ข
Kafka Connect ๋Š” Kafka Broker ์™€ ๋ณ„๋„์˜ ์„œ๋น„์Šค๋กœ ์šด์˜๋œ๋‹ค
โ€ข
๊ธฐ๋ณธ์ ์œผ๋กœ ํ•œ ํ…Œ์ด๋ธ”์˜ ๋ณ€๊ฒฝ ์‚ฌํ•ญ์€ ํ•˜๋‚˜์˜ ํ† ํ”ฝ์œผ๋กœ ์ „๋‹ฌ๋œ๋‹ค
โ€ข
MySQL ์˜ ๊ฒฝ์šฐ binlog ์— ์ ‘๊ทผํ•˜์—ฌ ๋ฐ์ดํ„ฐ๋ฅผ ๊ฐ€์ ธ์˜จ๋‹ค
โ€ข
PostgreSQL ์˜ ๊ฒฝ์šฐ Logical Replication Stream ์—์„œ ๋ฐ์ดํ„ฐ๋ฅผ ๊ฐ€์ ธ์˜จ๋‹ค
โ€ข
Source Connector ์—์„œ ๊ฐ€์ ธ์˜จ ์ •๋ณด๋“ค์„ Sink Connector ๋ฅผ ํ†ตํ•ด์„œ ElasticSearch, Redis, Data Warehouse ๋“ฑ์— ๋ฐ˜์˜ํ•  ์ˆ˜ ์žˆ๋‹ค.

Debezium Server

โ€ข
Kafka Connector ๋ฅผ ์‚ฌ์šฉํ•˜์ง€ ์•Š๊ณ  Debezium Server ๋ฅผ ์ด์šฉํ•ด์„œ ์‚ฌ์šฉํ•  ์ˆ˜ ์žˆ๋‹ค.
โ€ข
Debezium Server ๋Š” Source Connector ์ค‘ ํ•˜๋‚˜๋ฅผ ์‚ฌ์šฉํ•˜์—ฌ Source Database ์˜ ๋ณ€๊ฒฝ ์‚ฌํ•ญ์„ ์บก์ฒ˜ํ•˜๋„๋ก ๊ตฌ์„ฑํ•œ๋‹ค.
โ€ข
๋ณ€๊ฒฝ๋œ ์ด๋ฒคํŠธ๋Š” JSON, Apache Avro ์™€ ๊ฐ™์€ ํ˜•์‹์œผ๋กœ ์ง๋ ฌํ™” ํ•  ์ˆ˜ ์žˆ์œผ๋ฉฐ Kinesis, Google Pub/Sub, Redis ๋“ฑ์— ๋ฉ”์‹œ์ง€๋ฅผ ์ „๋‹ฌํ•  ์ˆ˜ ์žˆ๋‹ค

Embedded Engine

โ€ข
Kafka Connect ๋ฅผ ์‚ฌ์šฉํ•˜์ง€ ์•Š๊ณ  Embedded Engine ์„ ์‚ฌ์šฉํ•˜์—ฌ ์ž๋ฐ” ์• ํ”Œ๋ฆฌ์ผ€์ด์…˜ ๋ผ์ด๋ธŒ๋Ÿฌ๋ฆฌ๋กœ๋„ ์‚ฌ์šฉํ•  ์ˆ˜ ์žˆ๋‹ค.
โ—ฆ
Embedded Engine ์€ ๋ณ€๊ฒฝ ์ด๋ฒคํŠธ๋ฅผ ์• ํ”Œ๋ฆฌ์ผ€์ด์…˜์—์„œ ๋ฐ”๋กœ Consuming ํ•˜๊ฑฐ๋‚˜ ๋ณ€๊ฒฝ ๋‚ด์—ญ์„ ๋ณ„๋„์˜ ๋ฉ”์‹œ์ง€ ๋ธŒ๋กœ์ปค(e.g Amazon Kinesis) ์— ์ „๋‹ฌํ•  ๋•Œ ์œ ์šฉํ•˜๋‹ค.

๋กœ๊ทธ ๊ธฐ๋ฐ˜์˜ CDC ์˜ ์žฅ์ 

1.
๋ชจ๋“  ๋ฐ์ดํ„ฐ๋ฅผ ์บก์ฒ˜ํ•  ์ˆ˜ ์žˆ๋‹ค
โ€ข
๋กœ๊ทธ ๊ธฐ๋ฐ˜์œผ๋กœ ๋ฐ์ดํ„ฐ์˜ ๋ณ€๊ฒฝ์‚ฌํ•ญ์„ ๊ฐ€์ ธ์˜ค๊ฒŒ๋˜๋ฉด ์• ํ”Œ๋ฆฌ์ผ€์ด์…˜์—์„œ๋Š” ์ •ํ™•ํ•œ ์ˆœ์„œ๋Œ€๋กœ ๋ฐ์ดํ„ฐ์˜ ๋ณ€ํ™”๋ฅผ ๊ฐ์ง€ํ•  ์ˆ˜ ์žˆ๋‹ค.
โ€ข
ํด๋ง ๊ธฐ๋ฐ˜์˜ ๋ฐ์ดํ„ฐ ๋ณ€๊ฒฝ ์‚ฌํ•ญ์„ ๊ฐ€์ ธ์˜ค๊ฒŒ ๋  ๊ฒฝ์šฐ ์ตœ์‹  ๋ฐ์ดํ„ฐ๋งŒ ๊ฐ€์ ธ์˜ค๊ฒŒ ๋˜๋ฉฐ, ๋‹ค์šดํƒ€์ž„์ด ๋ฐœ์ƒ ํ–ˆ์„ ๋•Œ๋Š” ๋ฐ์ดํ„ฐ๋ฅผ ๋ชป๊ฐ€์ ธ์˜จ๋‹ค.
2.
CPU ๋ถ€ํ•˜๊ฐ€ ์ ๊ณ , ์ง€์—ฐ ์‹œ๊ฐ„์ด ์ ๋‹ค
โ€ข
์ฟผ๋ฆฌ๋ฅผ ์ฃผ๊ธฐ์ ์œผ๋กœ ์š”์ฒญํ•˜์ง€ ์•Š์•„๋„ ๋˜๊ธฐ ๋•Œ๋ฌธ์— CPU๋ฅผ ๊ฑฐ์˜ ์‚ฌ์šฉํ•˜์ง€ ํ•˜์ง€ ์•Š์œผ๋ฉฐ, ์‹ค์‹œ๊ฐ„์— ๊ฐ€๊นŒ์šด ๋ฐ์ดํ„ฐ ๋ณ€๊ฒฝ์— ๋Œ€์‘ํ•  ์ˆ˜ ์žˆ๋‹ค.
3.
๋ฐ์ดํ„ฐ ๋ชจ๋ธ์— ์˜ํ–ฅ์ด ์—†๋‹ค
โ€ข
ํด๋ง ๊ธฐ๋ฐ˜์˜ CDC ๋Š” ๋งˆ์ง€๋ง‰ ํด๋ง ์ดํ›„ ๋ณ€๊ฒฝ๋œ ๋ ˆ์ฝ”๋“œ๋ฅผ ์‹๋ณ„ํ•˜๊ธฐ ์œ„ํ•ด LAST_UPDATE_TIMESTAMP ์™€ ๊ฐ™์€ ์ •๋ณด๋“ค์„ ์ถ”๊ฐ€ํ•ด์•ผ๋˜์ง€๋งŒ, ๋กœ๊ทธ ๊ธฐ๋ฐ˜์˜ CDC ๋Š” ๊ทธ๋Ÿฐ๊ฒŒ ํ•„์š” ์—†๋‹ค.
4.
์‚ญ์ œ๋œ ๋ฐ์ดํ„ฐ๋ฅผ ์บก์ฒ˜ํ•  ์ˆ˜ ์žˆ๋‹ค.
5.
์ด์ „ ๋ ˆ์ฝ”๋“œ์˜ ์ƒํƒœ์™€ ๋ฉ”ํƒ€ ๋ฐ์ดํ„ฐ๋ฅผ ์บก์ฒ˜ํ•  ์ˆ˜ ์žˆ๋‹ค
โ€ข
๋ฐ์ดํ„ฐ๊ฐ€ ์—…๋ฐ์ดํŠธ ๋˜์—ˆ์„ ๋•Œ, ์ด์ „ ์ƒํƒœ์˜ ๋ ˆ์ฝ”๋“œ ์ƒํƒœ๋ฅผ ๊ฐ€์ ธ์˜ฌ ์ˆ˜ ์žˆ๋‹ค
โ€ข
๋กœ๊ทธ ๊ธฐ๋ฐ˜ CDC ๋Š” Schema Change Stream ์„ ์ œ๊ณตํ•˜๊ณ  Transaction ID ๋˜๋Š” ์‚ฌ์šฉ์ž๊ฐ€ ํŠน์ • ๋ณ€๊ฒฝ์‚ฌํ•ญ์„ ์ ์šฉํ•œ ๊ฒƒ๊ณผ ๊ฐ™์€ ๋ฉ”ํƒ€๋ฐ์ดํ„ฐ๋ฅผ ๊ฐ€์ ธ์˜ฌ ์ˆ˜ ์žˆ๋‹ค.

์‚ฌ์šฉ ๋ฐฉ๋ฒ•

MSK Connector

โ€ข
MSK Connector ๋Š” Kafka Connect ๋ฅผ ์ข€ ๋” ์‰ฝ๊ฒŒ ์‚ฌ์šฉํ•  ์ˆ˜ ์žˆ๋Š” ๊ธฐ๋Šฅ์„ ์ œ๊ณตํ•œ๋‹ค.
โ€ข
21.02.18 ๊ธฐ์ค€์œผ๋กœ Kafka Connect 2.7.1 ๋ฒ„์ „์„ ์‚ฌ์šฉํ•˜๊ณ  ์žˆ๋‹ค

์ œ๊ณตํ•˜๋Š” ๊ธฐ๋Šฅ

โ€ข
์ปค๋„ฅํ„ฐ์˜ ์ „์†ก ์ƒํƒœ ๋ชจ๋‹ˆํ„ฐ๋ง
โ€ข
ํ•˜๋“œ์›จ์–ด ํŒจ์น˜
โ€ข
์ฒ˜๋Ÿ‰ ๋ณ€ํ™”์— ๋งž์ถ˜ ์˜คํ†  ์Šค์ผ€์ผ๋ง
โ€ข
Kafka Connect ์™€ ์™„๋ฒฝํ•œ ํ˜ธํ™˜ ์ œ๊ณต
โ€ข
MSK ์™ธ์— On-Demand Kafka Cluster ์ง€์›
โ€ข
๊ทธ ์™ธ ์‚ฌ์šฉํ•˜๊ธฐ ์†์‰ฝ๊ฒŒ ๋งŒ๋“ค์–ด์ฃผ๋Š” ์™„์ „ ๊ด€๋ฆฌํ˜• ์„œ๋น„์Šค

MSK Connector ์•„ํ‚คํ…์ฒ˜

์‚ฌ์šฉ ๋ฐฉ๋ฒ•

Amazon MSK Connect - Apache Kafka ํด๋Ÿฌ์Šคํ„ฐ๋กœ ๋ฐ์ดํ„ฐ ์ „๋‹ฌ ์„œ๋น„์Šค ์ถœ์‹œ | Amazon Web Services
Apache Kafka๋Š” ์‹ค์‹œ๊ฐ„ ์ŠคํŠธ๋ฆฌ๋ฐ ๋ฐ์ดํ„ฐ ํŒŒ์ดํ”„๋ผ์ธ ๋ฐ ์• ํ”Œ๋ฆฌ์ผ€์ด์…˜ ๊ตฌ์ถ•์„ ์œ„ํ•œ ์˜คํ”ˆ ์†Œ์Šค ํ”Œ๋žซํผ์ž…๋‹ˆ๋‹ค. re:Invent 2018์—์„œ AWS๋Š” ์ŠคํŠธ๋ฆฌ๋ฐ ๋ฐ์ดํ„ฐ์˜ ํ”„๋กœ์„ธ์‹ฑ์„ ์œ„ํ•ด Apache Kafka๋ฅผ ์‚ฌ์šฉํ•˜๋Š” ์• ํ”Œ๋ฆฌ์ผ€์ด์…˜์„ ์‰ฝ๊ฒŒ ๊ตฌ์ถ• ๋ฐ ์‹คํ–‰ํ•  ์ˆ˜ ์žˆ๊ฒŒ ํ•ด ์ฃผ๋Š” ์™„์ „๊ด€๋ฆฌํ˜• ์„œ๋น„์Šค์ธ Amazon Managed Streaming for Apache Kafka๋ฅผ ๋ฐœํ‘œํ–ˆ์Šต๋‹ˆ๋‹ค. Apache Kafka๋ฅผ ์‚ฌ์šฉํ•˜๋ฉด IoT ๋””๋ฐ”์ด์Šค, ๋ฐ์ดํ„ฐ๋ฒ ์ด์Šค ๋ณ€๊ฒฝ ์ด๋ฒคํŠธ ๋ฐ ์›น ์‚ฌ์ดํŠธ ํด๋ฆญ์ŠคํŠธ๋ฆผ๊ณผ ๊ฐ™์€ ์†Œ์Šค๋กœ๋ถ€ํ„ฐ [...]