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 ๋ก ๋์ธ ์ ์๋ค
โข
โข
Guide
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
โข
โข
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 ์ง์
โข
๊ทธ ์ธ ์ฌ์ฉํ๊ธฐ ์์ฝ๊ฒ ๋ง๋ค์ด์ฃผ๋ ์์ ๊ด๋ฆฌํ ์๋น์ค







