๋ฐ์ดํฐ ์์ง๋์ด๋ฅผ ์ํ Apache Kafka: ์คํธ๋ฆฌ๋ฐ, ํํฐ์ , ๋ฉด์ ์ง๋ฌธ
๋ฐ์ดํฐ ์์ง๋์ด๋ฅผ ์ํ Apache Kafka ์ฌ์ธต ๋ถ์. Kafka 4.x์ KRaft๋ฅผ ํ์ฉํ ์คํธ๋ฆฌ๋ฐ ์ํคํ ์ฒ, ํํฐ์ ์ ๋ต, ์ปจ์๋จธ ๊ทธ๋ฃน, ๊ธฐ์ ๋ฉด์ ๋น์ถ ์ง๋ฌธ์ ์ค์ ์ฝ๋ ์์ ์ ํจ๊ป ์ค๋ช ํฉ๋๋ค.

Apache Kafka๋ ํ๋ ๋ฐ์ดํฐ ์์ง๋์ด๋ง ์คํ์ ํต์ฌ์ ์๋ฆฌํ๋ฉฐ, ๋ชจ๋ ๊ท๋ชจ์ ์กฐ์ง์์ ๋งค์ผ ์์กฐ ๊ฑด์ ์ด๋ฒคํธ๋ฅผ ์ฒ๋ฆฌํ๊ณ ์์ต๋๋ค. Kafka 4.x๋ถํฐ๋ KRaft ๊ธฐ๋ฐ์ ์์ ์์จ ์ด์์ด ๊ตฌํ๋์ด ZooKeeper ์์กด์ฑ์ด ์์ ํ ์ ๊ฑฐ๋์์ต๋๋ค. ๋ฆฌ์ผํ์ ๋ฐ์ดํฐ ํ์ดํ๋ผ์ธ์ ์ ๊ณ ํ์ค์ผ๋ก์ ๋์ผํ ์์ ์ฑ์ ์ ์งํ๋ฉด์๋, ์ด์ ๋ณต์ก์ฑ์ด ํฌ๊ฒ ๊ฐ์ํ์ต๋๋ค.
Apache Kafka 4.0๋ถํฐ ZooKeeper๊ฐ ๋ ์ด์ ํ์ํ์ง ์์ต๋๋ค. KRaft(Kafka Raft)๊ฐ ๋ชจ๋ ๋ฉํ๋ฐ์ดํฐ ๊ด๋ฆฌ๋ฅผ ๋ค์ดํฐ๋ธ๋ก ์ฒ๋ฆฌํ์ฌ ์ด์ ๋ณต์ก์ฑ์ ์ค์ด๊ณ , ํ๋ก๋์ ๋ฐฐํฌ์์ ์ธํ๋ผ ๊ตฌ์ฑ ์์๋ฅผ ํ๋ ์์ ํ ์ ๊ฑฐํ์ต๋๋ค.
๋ฐ์ดํฐ ํ์ดํ๋ผ์ธ์ ์ํ Kafka ์คํธ๋ฆฌ๋ฐ ์ํคํ ์ฒ
Kafka๋ ๋ถ์ฐ ์ปค๋ฐ ๋ก๊ทธ๋ก ๋์ํฉ๋๋ค. ํ๋ก๋์๊ฐ ํ ํฝ ํํฐ์ ์ ๋์ ๋ ์ฝ๋๋ฅผ ์ถ๊ฐํ๊ณ , ์ปจ์๋จธ๊ฐ ํด๋น ๋ ์ฝ๋๋ฅผ ์์๋๋ก ์ฝ์ต๋๋ค. ์ด ์ถ๊ฐ ์ ์ฉ(append-only) ์ค๊ณ ๋๋ถ์ ๋ฆฌ์ผํ์ ์คํธ๋ฆฌ๋ฐ๊ณผ ์์ ์คํ์ ์์์ ๋ฐฐ์น ๋ฆฌํ๋ ์ด๊ฐ ๋ชจ๋ ๊ฐ๋ฅํ๋ฉฐ, ๋ ๊ฐ์ง ํจํด์ ๋ชจ๋ ํ์๋ก ํ๋ ๋ฐ์ดํฐ ์์ง๋์ด๋ง ์ํฌ๋ก๋์ ์ต์ ์ ๊ธฐ๋ฐ์ ์ ๊ณตํฉ๋๋ค.
Kafka ํด๋ฌ์คํฐ๋ ๋ธ๋ก์ปค(์๋ฒ), ํ ํฝ(๋ ผ๋ฆฌ์ ์ฑ๋), ํํฐ์ (ํ ํฝ์ ๋ฌผ๋ฆฌ์ ์ค๋)์ผ๋ก ๊ตฌ์ฑ๋ฉ๋๋ค. ๊ฐ ํํฐ์ ์ ์์๊ฐ ์ง์ ๋ ๋ถ๋ณ์ ๋ ์ฝ๋ ์ํ์ค์ ๋๋ค. ํ๋์ ๋ธ๋ก์ปค๊ฐ ํํฐ์ ๋ฆฌ๋ ์ญํ ์ ํ๋ฉฐ ๋ชจ๋ ์ฝ๊ธฐ์ ์ฐ๊ธฐ๋ฅผ ์ฒ๋ฆฌํฉ๋๋ค. ํ๋ก์ ๋ ํ๋ฆฌ์นด๋ ๋ด๊ฒฐํจ์ฑ์ ์ํด ์ฌ๋ณธ์ ์ ์งํฉ๋๋ค.
ํต์ฌ ์ํคํ ์ฒ ์์ฑ์ ๋ค์๊ณผ ๊ฐ์ต๋๋ค: ์์ ๋ณด์ฅ์ ํํฐ์ ๋ด์์๋ง ์ ํจํ๋ฉฐ, ํํฐ์ ๊ฐ์๋ ๋ณด์ฅ๋์ง ์์ต๋๋ค. ์ด ๋จ์ผ ์ ์ฝ ์กฐ๊ฑด์ด ๋ฐ์ดํฐ ํ์ดํ๋ผ์ธ์์ ๋๋ถ๋ถ์ Kafka ์ค๊ณ ๊ฒฐ์ ์ ์ข์ฐํฉ๋๋ค.
# docker-compose.yml โ Kafka 4.x KRaft cluster (no ZooKeeper)
services:
kafka-1:
image: apache/kafka:4.2.0
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_LOG_DIRS: /var/lib/kafka/data
KAFKA_NUM_PARTITIONS: 6
KAFKA_DEFAULT_REPLICATION_FACTOR: 3
KAFKA_MIN_INSYNC_REPLICAS: 2
ports:
- "9092:9092"์ด ๊ตฌ์ฑ์ ์ปจํธ๋กค๋ฌ ์ญํ ๋ ๊ฒธํ๋ KRaft ๋ชจ๋ ๋ธ๋ก์ปค๋ฅผ ์์ํฉ๋๋ค. KAFKA_CONTROLLER_QUORUM_VOTERS ์ค์ ์ ์ด์ ์ ZooKeeper๊ฐ ๊ด๋ฆฌํ๋ ๋ฆฌ๋ ์ ์ถ, ํ ํฝ ๋ฉํ๋ฐ์ดํฐ, ํด๋ฌ์คํฐ ๋ฉค๋ฒ์ญ ์ญํ ์ ๋์ฒดํฉ๋๋ค.
ํ์ดํ๋ผ์ธ ์ฑ๋ฅ์ ๊ฒฐ์ ํ๋ ํํฐ์ ์ ๋ต
ํํฐ์ ์ ๋ณ๋ ฌ ์ฒ๋ฆฌ ์์ค, ์์ ๋ณด์ฅ, ์ฒ๋ฆฌ๋์ ๊ฒฐ์ ํฉ๋๋ค. ์ ์ ํ ํํฐ์ ์์ ํค ์ ๋ต์ ์ ํ์ด ํ์ดํ๋ผ์ธ์ ์์ ์ฑ๊ณผ ์ปจ์๋จธ ํ์ฅ์ฑ์ ์ง์ ์ ์ธ ์ํฅ์ ๋ฏธ์นฉ๋๋ค.
ํํฐ์ ์ ๊ฐ์ด๋๋ผ์ธ:
- ํ ํฝ์ ๋์์ ์ฒ๋ฆฌํ ๊ฒ์ผ๋ก ์์๋๋ ์ปจ์๋จธ ์๋ก ์์ํฉ๋๋ค
- ๊ฐ ํํฐ์ ์ ์ต์ ํ๋์จ์ด์์ ์ฝ 10MB/s์ ์ฐ๊ธฐ ์ฒ๋ฆฌ๋์ ์ ์งํ ์ ์์ต๋๋ค
- ํํฐ์ ์๊ฐ ๋ง์์ง๋ฉด ์๋ํฌ์๋ ์ง์ฐ ์๊ฐ์ด ์ฝ๊ฐ ์ฆ๊ฐํฉ๋๋ค(๋ฆฌ๋ ์ ์ถ๊ณผ ํ์ผ ํธ๋ค ์ฆ๊ฐ)
- Kafka 4.x๋ KRaft์ ๋ฉํ๋ฐ์ดํฐ ๊ฐ์ ๋๋ถ์ ๋ธ๋ก์ปค๋น ์์ฒ ๊ฐ์ ํํฐ์ ์ ํจ์จ์ ์ผ๋ก ์ฒ๋ฆฌํฉ๋๋ค
ํํฐ์ ํค ์ ํ์ ๋ฐ๋ผ ์ด๋ค ๋ ์ฝ๋๊ฐ ์ด๋ค ํํฐ์ ์ ๋ฐฐ์น๋ ์ง ๊ฒฐ์ ๋ฉ๋๋ค. ๊ฐ์ ํค๋ฅผ ๊ฐ์ง ๋ ์ฝ๋๋ ํญ์ ๊ฐ์ ํํฐ์ ์ผ๋ก ์ ์ก๋์ด, ํด๋น ํค์ ๋ํ ์์๊ฐ ๋ณด์กด๋ฉ๋๋ค.
# partition_strategy.py โ Choosing the right partition key
from confluent_kafka import Producer
import json
def create_producer():
return Producer({
'bootstrap.servers': 'kafka-1:9092',
'acks': 'all',
'enable.idempotence': True,
'max.in.flight.requests.per.connection': 5
})
def publish_order_event(producer, event):
# Key by customer_id: all events for one customer stay ordered
key = str(event['customer_id']).encode('utf-8')
value = json.dumps(event).encode('utf-8')
producer.produce(
topic='order-events',
key=key,
value=value,
headers=[('source', b'order-service'), ('version', b'2')]
)
producer.flush()
def publish_clickstream(producer, event):
# Key by session_id: ordering within a session matters
# NOT by user_id โ one user may have multiple sessions
key = str(event['session_id']).encode('utf-8')
value = json.dumps(event).encode('utf-8')
producer.produce(
topic='clickstream',
key=key,
value=value
)customer_id์ session_id ์ค ๋ฌด์์ ํํฐ์
ํค๋ก ์ฌ์ฉํ ์ง๋ ์๋ก ๋ค๋ฅธ ์์ ์๊ตฌ์ฌํญ์ ๋ฐ์ํฉ๋๋ค. ์ฃผ๋ฌธ ์ด๋ฒคํธ๋ ํธ๋์ญ์
์ผ๊ด์ฑ์ ์ ์งํ๊ธฐ ์ํด ๊ณ ๊ฐ ๋จ์์ ์์๊ฐ ํ์ํฉ๋๋ค. ํด๋ฆญ์คํธ๋ฆผ ์ด๋ฒคํธ๋ ์ฌ์ฉ์ ์ฌ์ ์ ์ ํํ๊ฒ ์ฌ๊ตฌ์ฑํ๊ธฐ ์ํด ์ธ์
๋จ์์ ์์๊ฐ ํ์ํฉ๋๋ค.
ํ๋์ ํค๊ฐ ๋ค๋ฅธ ํค๋ณด๋ค ์๋์ ์ผ๋ก ๋ง์ ๋ ์ฝ๋๋ฅผ ์์ฑํ๋ ๊ฒฝ์ฐ(์: ๋จ์ผ ์ํฐํ๋ผ์ด์ฆ ๊ณ ๊ฐ์ด ์ด๋ฒคํธ์ 80%๋ฅผ ์์ฑํ๋ ์ํฉ), ํด๋น ํํฐ์ ์ด ๋ณ๋ชฉ ์ง์ ์ด ๋ฉ๋๋ค. ํํฐ์ ๋ ๋ฉํธ๋ฆญ์ ๋ชจ๋ํฐ๋งํ๊ณ , ํธํฅ๋ ์ํฌ๋ก๋์๋ ๋ณตํฉ ํค๋ ์ปค์คํ ํํฐ์ ๋ ์ ์ฉ์ ๊ฒํ ํด์ผ ํฉ๋๋ค.
์ปจ์๋จธ ๊ทธ๋ฃน๊ณผ ์คํ์ ๊ด๋ฆฌ
์ปจ์๋จธ ๊ทธ๋ฃน์ ํตํด ํ ํฝ์ ๋ณ๋ ฌ ์ฒ๋ฆฌ๊ฐ ๊ฐ๋ฅํด์ง๋๋ค. Kafka๋ ๊ทธ๋ฃน ๋ด ๊ฐ ํํฐ์ ์ ์ ํํ ํ๋์ ์ปจ์๋จธ์ ํ ๋นํ๋ฏ๋ก, ์ต๋ ๋ณ๋ ฌ ์ฒ๋ฆฌ ์์ค์ ํํฐ์ ์์ ๋์ผํฉ๋๋ค.
์คํ์
๊ด๋ฆฌ๋ Exactly-once(์ ํํ ํ ๋ฒ)์ At-least-once(์ต์ ํ ๋ฒ) ์๋งจํฑ์ ์ ์ดํฉ๋๋ค. Kafka๋ ์ปค๋ฐ๋ ์คํ์
์ ๋ด๋ถ __consumer_offsets ํ ํฝ์ ์ ์ฅํฉ๋๋ค. ์ฒ๋ฆฌ ๋๋น ์คํ์
์ปค๋ฐ์ ํ์ด๋ฐ์ด ์ ๋ฌ ๋ณด์ฅ ์์ค์ ๊ฒฐ์ ํฉ๋๋ค.
# consumer_pipeline.py โ Consumer group with manual offset management
from confluent_kafka import Consumer, KafkaError
import json
def create_consumer(group_id: str):
return Consumer({
'bootstrap.servers': 'kafka-1:9092',
'group.id': group_id,
'auto.offset.reset': 'earliest',
'enable.auto.commit': False, # Manual commit for at-least-once
'max.poll.interval.ms': 300000,
'session.timeout.ms': 45000,
'isolation.level': 'read_committed' # Only read committed transactions
})
def process_batch(consumer: Consumer, batch_size: int = 500):
"""Process records in micro-batches for throughput."""
consumer.subscribe(['order-events'])
buffer = []
while True:
msg = consumer.poll(timeout=1.0)
if msg is None:
continue
if msg.error():
if msg.error().code() == KafkaError._PARTITION_EOF:
continue
raise Exception(msg.error())
record = json.loads(msg.value().decode('utf-8'))
buffer.append(record)
if len(buffer) >= batch_size:
# Process the full batch (write to warehouse, etc.)
write_to_warehouse(buffer)
# Commit AFTER successful processing
consumer.commit(asynchronous=False)
buffer.clear()
def write_to_warehouse(records: list):
# Insert batch into data warehouse (BigQuery, Snowflake, etc.)
passenable.auto.commit๋ฅผ False๋ก ์ค์ ํ๊ณ write_to_warehouse ์ฑ๊ณต ํ์๋ง ์ปค๋ฐํจ์ผ๋ก์จ At-least-once ์ ๋ฌ์ด ๋ณด์ฅ๋ฉ๋๋ค. ์ปจ์๋จธ๊ฐ ์ฒ๋ฆฌ ํ ์ปค๋ฐ ์ ์ ํฌ๋์๋๋ฉด, ์ฌ์์ ์ ๋ ์ฝ๋๊ฐ ์ฌ์ ๋ฌ๋ฉ๋๋ค. ์จ์ดํ์ฐ์ค์ ๋ฐ์ดํฐ๋ฅผ ๊ณต๊ธํ๋ ํ์ดํ๋ผ์ธ์์๋ ์ด๊ฒ์ด ์ผ๋ฐ์ ์ผ๋ก ์ฌ๋ฐ๋ฅธ ํธ๋ ์ด๋์คํ์
๋๋ค. ์จ์ดํ์ฐ์ค์ ๋ํ ๋ฉฑ๋ฑ ์ฐ๊ธฐ๊ฐ ์ค๋ณต์ ์ฒ๋ฆฌํฉ๋๋ค.
๋ฐ์ดํฐ ํ์ดํ๋ผ์ธ ํตํฉ์ ์ํ Kafka Connect
Kafka Connect๋ ์ปค์คํ ํ๋ก๋์/์ปจ์๋จธ ์ฝ๋๋ฅผ ์์ฑํ์ง ์๊ณ ๋ Kafka์ ์ธ๋ถ ์์คํ ๊ฐ์ ๋ฐ์ดํฐ ์ด๋์ ๊ฐ๋ฅํ๊ฒ ํ๋ ํ์ฅ ๊ฐ๋ฅํ ํ๋ ์์ํฌ์ ๋๋ค. ์์ค ์ปค๋ฅํฐ๊ฐ Kafka๋ก ๋ฐ์ดํฐ๋ฅผ ์์งํ๊ณ , ์ฑํฌ ์ปค๋ฅํฐ๊ฐ ๋ฐ์ดํฐ๋ฅผ ์ธ๋ถ๋ก ์ ์กํฉ๋๋ค.
{
"name": "postgres-cdc-source",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "postgres-primary",
"database.port": "5432",
"database.user": "replication_user",
"database.dbname": "production",
"topic.prefix": "cdc",
"table.include.list": "public.orders,public.customers,public.products",
"slot.name": "debezium_slot",
"plugin.name": "pgoutput",
"publication.name": "dbz_publication",
"snapshot.mode": "initial",
"transforms": "route",
"transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "cdc\\.public\\.(.*)",
"transforms.route.replacement": "warehouse.$1"
}
}์ด Debezium CDC ์ปค๋ฅํฐ๋ PostgreSQL ํ
์ด๋ธ์์ ๋ฐ์ํ๋ ๋ชจ๋ INSERT, UPDATE, DELETE๋ฅผ ์บก์ฒํ์ฌ Kafka ์ด๋ฒคํธ๋ก ๋ฐํํฉ๋๋ค. RegexRouter ํธ๋์คํผ์ด ํ ํฝ ์ด๋ฆ์ cdc.public.orders์์ warehouse.orders๋ก ๋ณ๊ฒฝํ์ฌ ํ์ดํ๋ผ์ธ ๋ค์์คํ์ด์ค๋ฅผ ์ ๋ฆฌํฉ๋๋ค. BigQuery๋ Snowflake์ ์ฑํฌ ์ปค๋ฅํฐ์ ๊ฒฐํฉํ๋ฉด, ์ปค์คํ
์ฝ๋ ์์ด ์์ ์๋ํ๋ CDC ํ์ดํ๋ผ์ธ์ ๊ตฌ์ถํ ์ ์์ต๋๋ค.
Data Engineering ๋ฉด์ ์ค๋น๊ฐ ๋์ จ๋์?
์ธํฐ๋ํฐ๋ธ ์๋ฎฌ๋ ์ดํฐ, flashcards, ๊ธฐ์ ํ ์คํธ๋ก ์ฐ์ตํ์ธ์.
Kafka ํ์ดํ๋ผ์ธ์ Exactly-Once ์๋งจํฑ
Kafka๋ ๋ฉฑ๋ฑ ํ๋ก๋์์ ํธ๋์ญ์ ๋ API๋ฅผ ํตํด Exactly-once ์๋งจํฑ(EOS)์ ์ง์ํฉ๋๋ค. ๋ฐ์ดํฐ ์์ง๋์ด๋ง ํ์ดํ๋ผ์ธ์์ EOS๋ ์ค๋ณต ์ ๊ฑฐ ๋ก์ง ์์ด๋ ๋ค์ด์คํธ๋ฆผ ์์คํ ์ ์ค๋ณต ๋ ์ฝ๋๋ฅผ ๋ฐฉ์งํฉ๋๋ค.
EOS๋ฅผ ๊ตฌํํ๋ 3๊ฐ์ง ๊ตฌ์ฑ ์์:
- ๋ฉฑ๋ฑ ํ๋ก๋์: ๊ฐ ํ๋ก๋์๋ ๊ณ ์ ํ Producer ID๋ฅผ ๋ถ์ฌ๋ฐ์ต๋๋ค. ๋ธ๋ก์ปค๋ ์ํ์ค ๋ฒํธ๋ฅผ ์ฌ์ฉํ์ฌ ์ฌ์๋์ ์ค๋ณต์ ์ ๊ฑฐํ๋ฏ๋ก, ๋คํธ์ํฌ ์ฌ์๋๊ฐ ์ค๋ณต์ ์์ฑํ์ง ์์ต๋๋ค.
- ํธ๋์ญ์ : ํ๋ก๋์๋ ์ฌ๋ฌ ํํฐ์ ์ ๋ํ ์ฐ๊ธฐ์ ์คํ์ ์ปค๋ฐ์ ๋จ์ผ ํธ๋์ญ์ ์ผ๋ก ์์์ ์ผ๋ก ์ํํ ์ ์์ต๋๋ค. ๋ชจ๋ ์ฐ๊ธฐ๊ฐ ์ฑ๊ณตํ๊ฑฐ๋ ๋ชจ๋ ์คํจํฉ๋๋ค.
read_committed๊ฒฉ๋ฆฌ ์์ค:isolation.level=read_committed๋ก ์ค์ ๋ ์ปจ์๋จธ๋ ์ปค๋ฐ๋ ํธ๋์ญ์ ์ ๋ ์ฝ๋๋ง ์ฝ์ต๋๋ค.
# exactly_once_pipeline.py โ Transactional consume-transform-produce
from confluent_kafka import Consumer, Producer
import json
def transactional_pipeline():
consumer = Consumer({
'bootstrap.servers': 'kafka-1:9092',
'group.id': 'etl-transformer',
'enable.auto.commit': False,
'isolation.level': 'read_committed'
})
producer = Producer({
'bootstrap.servers': 'kafka-1:9092',
'transactional.id': 'etl-transformer-001',
'acks': 'all',
'enable.idempotence': True
})
producer.init_transactions()
consumer.subscribe(['raw-events'])
while True:
msg = consumer.poll(1.0)
if msg is None or msg.error():
continue
# Begin atomic transaction
producer.begin_transaction()
try:
raw = json.loads(msg.value())
enriched = transform(raw)
# Write transformed record to output topic
producer.produce(
'enriched-events',
key=msg.key(),
value=json.dumps(enriched).encode('utf-8')
)
# Commit consumer offset within the same transaction
producer.send_offsets_to_transaction(
consumer.position(consumer.assignment()),
consumer.consumer_group_metadata()
)
producer.commit_transaction()
except Exception:
producer.abort_transaction()
def transform(record: dict) -> dict:
# Apply business logic transformations
record['processed_at'] = '2026-04-20T00:00:00Z'
record['pipeline_version'] = '2.1'
return recordtransactional.id ์ค์ ์ ํตํด ๋ธ๋ก์ปค๋ ์ข๋น ํ๋ก๋์๋ฅผ ์ฐจ๋จํ ์ ์์ต๋๋ค. ์ปจ์๋จธ๊ฐ ํฌ๋์๋๊ณ ๋์ผํ transactional.id๋ก ์ ์ธ์คํด์ค๊ฐ ์์๋๋ฉด, ๋ธ๋ก์ปค๋ ์ด์ ์ธ์คํด์ค์ ๋ณด๋ฅ ์ค์ธ ํธ๋์ญ์
์ ์ค๋จํ์ฌ ์ถ๋ ฅ ์ค๋ณต์ ๋ฐฉ์งํฉ๋๋ค.
Share Groups: Kafka 4.2์ ํ ์๋งจํฑ
Kafka 4.2์์๋ ํ๋ก๋์ ์์ค์ Share Groups๊ฐ ๋์ ๋์ด, Kafka์ ์ ํต์ ์ธ ๋ฉ์์ง ํ ์๋งจํฑ์ด ์ถ๊ฐ๋์์ต๋๋ค. ๊ฐ ํํฐ์ ์ด ๊ทธ๋ฃน ๋ด ์ ํํ ํ๋์ ์ปจ์๋จธ์๋ง ํ ๋น๋๋ ์ปจ์๋จธ ๊ทธ๋ฃน๊ณผ ๋ฌ๋ฆฌ, Share Groups์์๋ ์ฌ๋ฌ ์ปจ์๋จธ๊ฐ ๋์ผํ ํํฐ์ ์ ๋ ์ฝ๋๋ฅผ ๋ ๋ฆฝ์ ์ผ๋ก ์ฒ๋ฆฌํ ์ ์์ต๋๋ค.
| ํน์ฑ | ์ปจ์๋จธ ๊ทธ๋ฃน | Share Groups |
|---|---|---|
| ํํฐ์ ํ ๋น | ๋ฐฐํ์ (ํํฐ์ ๋น ์ปจ์๋จธ 1๊ฐ) | ๊ณต์ (ํํฐ์ ๋น ๋ณต์ ์ปจ์๋จธ) |
| ์์ ๋ณด์ฅ | ํํฐ์ ๋ด ์์ ๋ณด์กด | ์์ ๋ณด์ฅ ์์ |
| ์ฌ์ฉ ์ฌ๋ก | ์์ ๋ณด์ฅ ์ด๋ฒคํธ ์คํธ๋ฆผ | ์์ ๋ถ๋ฐฐ, ์์ ํ |
| ํ์ธ ์๋ต | ์คํ์ ๊ธฐ๋ฐ (์์น ์ปค๋ฐ) | ๋ ์ฝ๋ ๋จ์ (๊ฐ๋ณ ํ์ธ ์๋ต) |
| ์ต๋ ๋ณ๋ ฌ ์ฒ๋ฆฌ | ํํฐ์ ์๋ก ์ ํ | ์ปจ์๋จธ ์๋ก ์ ํ |
Share Groups๋ ์ค๋ซ๋์ ์กด์ฌํ๋ ํ๊ณ, ์ฆ ํํฐ์ ์๋ฅผ ์ด๊ณผํ๋ ์ปจ์๋จธ ํ์ฅ ๋ฌธ์ ๋ฅผ ํด๊ฒฐํฉ๋๋ค. ์๋ฆผ ์ ์ก์ด๋ ๋ฐฐ์น ์์ ๋ถ๋ฐฐ์ ๊ฐ์ด ์์๊ฐ ํ์ ์๋ ์ํฌ๋ก๋์์๋ Kafka ์ธ์ ๋ณ๋์ ์ธ๋ถ ํ ์์คํ ์ด ํ์ํ์ง ์๊ฒ ๋ฉ๋๋ค.
Share Groups๋ ์์ ๋ถ๋ฐฐํ ์ํฌ๋ก๋์ ์ ํฉํฉ๋๋ค: ์ด๋ฉ์ผ ๋ฐ์ก, ์ ๋ก๋ ์ฒ๋ฆฌ, ML ์ถ๋ก ์คํ ๋ฑ. ์ด๋ฒคํธ ์์ฑ, CDC, ๋๋ ์์ ๋ณด์ฅ ์ฒ๋ฆฌ๊ฐ ํ์ํ ํ์ดํ๋ผ์ธ์์๋ ์ปจ์๋จธ ๊ทธ๋ฃน์ด ์ฌ์ ํ ์ฌ๋ฐ๋ฅธ ์ ํ์ ๋๋ค.
๋ฐ์ดํฐ ์์ง๋์ด๋ฅผ ์ํ Kafka ๋ฉด์ ์ง๋ฌธ
๋ฐ์ดํฐ ์์ง๋์ด๋ง ์ง๋ฌด์ ๊ธฐ์ ๋ฉด์ ์์๋ Kafka ๊ด๋ จ ์ง์์ด ์์ฃผ ํ๊ฐ๋ฉ๋๋ค. ๋ค์ ์ง๋ฌธ๋ค์ ๊ฐ์ฅ ์ผ๋ฐ์ ์ผ๋ก ํ๊ฐ๋๋ ๊ฐ๋ ์ ๋ค๋ฃจ๊ณ ์์ต๋๋ค.
Q: Kafka๋ ๋ฉ์์ง ์์๋ฅผ ์ด๋ป๊ฒ ๋ณด์ฅํฉ๋๊น? ์์๋ ๋จ์ผ ํํฐ์ ๋ด์์๋ง ๋ณด์ฅ๋ฉ๋๋ค. ๋์ผํ ํํฐ์ ํค๋ฅผ ๊ฐ์ง ๋ชจ๋ ๋ ์ฝ๋๋ ๊ฐ์ ํํฐ์ ์ผ๋ก ์ ์ก๋์ด ์์๋๋ก ์ถ๊ฐ๋ฉ๋๋ค. ํํฐ์ ๊ฐ์๋ ์์๊ฐ ์กด์ฌํ์ง ์์ต๋๋ค. ์ ์ ํ ํํฐ์ ํค ์ ํ์ด Kafka ํ์ดํ๋ผ์ธ์์ ์์๋ฅผ ์ ์ดํ๋ ์ฃผ์ ๋ฉ์ปค๋์ฆ์ ๋๋ค.
Q: ๊ทธ๋ฃน ๋ด ์ปจ์๋จธ์ ์ฅ์ ๊ฐ ๋ฐ์ํ๋ฉด ์ด๋ป๊ฒ ๋ฉ๋๊น?
๊ทธ๋ฃน ์ฝ๋๋ค์ดํฐ๊ฐ ํํธ๋นํธ ๋๋ฝ์ ํตํด ์ฅ์ ๋ฅผ ๊ฐ์งํฉ๋๋ค(session.timeout.ms๋ก ์ ์ด). ๋ฆฌ๋ฐธ๋ฐ์ค๊ฐ ํธ๋ฆฌ๊ฑฐ๋์ด, ์ฅ์ ๊ฐ ๋ฐ์ํ ์ปจ์๋จธ์ ํํฐ์
์ด ๋๋จธ์ง ๊ทธ๋ฃน ๋ฉค๋ฒ์๊ฒ ์ฌํ ๋น๋ฉ๋๋ค. ๋ฆฌ๋ฐธ๋ฐ์ค ์ค์๋ ์๋น๊ฐ ์ผ์์ ์ผ๋ก ์ค๋จ๋ฉ๋๋ค. Kafka 4.x์์ ๊ธฐ๋ณธ๊ฐ์ธ ํ๋ ฅ์ ์คํฐํค ๋ฆฌ๋ฐธ๋ฐ์ฑ์ ์ํฅ๋ฐ๋ ํํฐ์
๋ง ์ฌํ ๋นํ์ฌ ์ค๋จ์ ์ต์ํํฉ๋๋ค.
Q: acks=1๊ณผ acks=all์ ์ฐจ์ด๋ฅผ ์ค๋ช
ํด ์ฃผ์ญ์์ค.
acks=1์์๋ ๋ฆฌ๋ ๋ ํ๋ฆฌ์นด๊ฐ ๋ ์ฝ๋๋ฅผ ์๊ตฌ ์ ์ฅํ ํ ๋ธ๋ก์ปค๊ฐ ์ฐ๊ธฐ๋ฅผ ํ์ธ ์๋ตํฉ๋๋ค. acks=all์์๋ ๋ชจ๋ In-Sync ๋ ํ๋ฆฌ์นด(ISR)๊ฐ ์ฐ๊ธฐ๋ฅผ ํ์ธํ ๋๊น์ง ๋๊ธฐํฉ๋๋ค. acks=all์ min.insync.replicas=2์ ๊ฒฐํฉํ๋ฉด, ๋จ์ผ ๋ธ๋ก์ปค ์ฅ์ ์์๋ ๋ฐ์ดํฐ ์์ค์ด ์์์ ๋ณด์ฅํฉ๋๋ค. ์ง์ฐ ์๊ฐ์ ์ฝ๊ฐ ์ฆ๊ฐํฉ๋๋ค.
Q: Kafka๋ฅผ ์ฌ์ฉํ CDC ํ์ดํ๋ผ์ธ์ ์ด๋ป๊ฒ ์ค๊ณํฉ๋๊น?
Debezium(PostgreSQL, MySQL์ฉ) ๋๋ ๋ค์ดํฐ๋ธ CDC ์ปค๋ฅํฐ๋ฅผ ์ฌ์ฉํ์ฌ ์์ค ๋ฐ์ดํฐ๋ฒ ์ด์ค์ ๋ณ๊ฒฝ ์ฌํญ์ ์บก์ฒํฉ๋๋ค. ๋ณ๊ฒฝ ์ด๋ฒคํธ๋ฅผ ๊ธฐ๋ณธ ํค๋ก ํํฐ์
๋๋ Kafka ํ ํฝ์ ๋ฐํํฉ๋๋ค. ์ฑํฌ ์ปค๋ฅํฐ(BigQuery Sink, S3 Sink)๋ฅผ ์ฌ์ฉํ์ฌ ๋ณ๊ฒฝ ์ฌํญ์ ๋ถ์์ฉ ์จ์ดํ์ฐ์ค์ ๋ก๋ํฉ๋๋ค. ์ปจ์๋จธ์ isolation.level=read_committed๋ฅผ ์ค์ ํ์ฌ ์ปค๋ฐ๋์ง ์์ ๋ฐ์ดํฐ๋ฒ ์ด์ค ํธ๋์ญ์
์ ์ฝ๊ธฐ๋ฅผ ๋ฐฉ์งํฉ๋๋ค. ๋ณ๊ฒฝ ์ด๋ฒคํธ์ Avro ๋๋ Protobuf ์คํค๋ง ๊ด๋ฆฌ์๋ ์คํค๋ง ๋ ์ง์คํธ๋ฆฌ๋ฅผ ์ฌ์ฉํฉ๋๋ค.
Q: ISR์ด๋ ๋ฌด์์ด๋ฉฐ ์ ์ค์ํฉ๋๊น?
ISR(In-Sync Replicas)์ ํํฐ์
๋ฆฌ๋์ ์์ ํ ๋๊ธฐํ๋ ๋ ํ๋ฆฌ์นด์ ์งํฉ์
๋๋ค. ๋ฆฌ๋ ์ ์ถ ๋์์ ISR ๋ฉค๋ฒ๋ง ํด๋น๋ฉ๋๋ค. min.insync.replicas ์ค์ ์ ์ฐ๊ธฐ๊ฐ ์ปค๋ฐ๋ ๊ฒ์ผ๋ก ๊ฐ์ฃผ๋๊ธฐ ์ํด ํ์ธ ์๋ต์ด ํ์ํ ๋ ํ๋ฆฌ์นด ์๋ฅผ ์ ์ํฉ๋๋ค. ISR์ด min.insync.replicas ๋ฏธ๋ง์ผ๋ก ์ค์ด๋ค๋ฉด, ๋ธ๋ก์ปค๋ ๋ฐ์ดํฐ ์์ค ์ํ์ ํผํ๊ธฐ ์ํด ์ฐ๊ธฐ๋ฅผ ๊ฑฐ๋ถํฉ๋๋ค.
์ถ๊ฐ์ ์ธ ๋ฐ์ดํฐ ์์ง๋์ด๋ง ๋ฉด์ ์ค๋น๋ฅผ ์ํด, ๋ฐ์ดํฐ ์์ง๋์ด๋ง ๋ฉด์ ์ง๋ฌธ ๋ชจ์์์ ETL/ELT ํ์ดํ๋ผ์ธ ํจํด๊ณผ ๋ฐ์ดํฐ ๋ชจ๋ธ๋ง ๋ฑ ๋ ๋์ ์ฃผ์ ๋ฅผ ๋ค๋ฃจ๊ณ ์์ต๋๋ค.
์ฐ์ต์ ์์ํ์ธ์!
๋ฉด์ ์๋ฎฌ๋ ์ดํฐ์ ๊ธฐ์ ํ ์คํธ๋ก ์ง์์ ํ ์คํธํ์ธ์.
๊ฒฐ๋ก
- Kafka 4.x์ KRaft๋ก ZooKeeper๊ฐ ์์ ํ ์ ๊ฑฐ๋์ด, ๋ฐ์ดํฐ ์์ง๋์ด๋ง ํ์ ์ด์ ์ค๋ฒํค๋๊ฐ ๊ฐ์ํฉ๋๋ค
- ํํฐ์ ํค ์ ํ์ด ์์ ๋ณด์ฅ์ ๊ฒฐ์ ํฉ๋๋ค. ํธ์๊ฐ ์๋ ๋ค์ด์คํธ๋ฆผ ์ฒ๋ฆฌ ์๊ตฌ์ฌํญ์ ๋ฐ๋ผ ํค๋ฅผ ์ ํํด์ผ ํฉ๋๋ค
enable.auto.commit=False๋ฅผ ํตํ ์๋ ์คํ์ ์ปค๋ฐ์ผ๋ก At-least-once ์ ๋ฌ์ด ๊ฐ๋ฅํ๋ฉฐ, ํธ๋์ญ์ ๋ API๋ consume-transform-produce ํ์ดํ๋ผ์ธ์ Exactly-once ์๋งจํฑ์ ์ ๊ณตํฉ๋๋ค- Kafka Connect์ Debezium์ ํ์ฉํ๋ฉด ์ปค์คํ ์ปจ์๋จธ ์ฝ๋ ์์ด ํ๋ก๋์ ์์ค์ CDC๋ฅผ ๊ตฌํํ ์ ์์ต๋๋ค
- Share Groups(Kafka 4.2)๋ ์์๊ฐ ํ์ ์๋ ์์ ๋ถ๋ฐฐ ์ํฌ๋ก๋์ ํ ์๋งจํฑ์ ์ถ๊ฐํฉ๋๋ค
- ์ปจ์๋จธ ๊ทธ๋ฃน ํ์ฅ์ ํํฐ์ ์์ ์ํด ์ ํ๋๋ฏ๋ก, ํผํฌ ์ ์ปจ์๋จธ ๋ณ๋ ฌ ์ฒ๋ฆฌ ์๊ตฌ์ ๊ธฐ๋ฐํ์ฌ ํํฐ์ ์๋ฅผ ๊ณํํด์ผ ํฉ๋๋ค
์ฐ์ต์ ์์ํ์ธ์!
๋ฉด์ ์๋ฎฌ๋ ์ดํฐ์ ๊ธฐ์ ํ ์คํธ๋ก ์ง์์ ํ ์คํธํ์ธ์.
Data Engineering ์ฝ๋์ ๋ฒ๊ทธ๋ฅผ ์ฐพ์ ์ ์๋์
์ค์ ์ฝ๋ ํ ์กฐ๊ฐ, ์จ์ ๋ฒ๊ทธ ํ๋, ํ๋ฃจ ํ ๋ฒ. ๊ณ์ ์์ด ๋ฐ๋ก ๋์ ํ ์ ์์ต๋๋ค.

์์ฑ์
Anthony Fillion-MailletSharpSkill ์ฐฝ์ ์
10๋ ์ด์ ํ์คํ ๊ฐ๋ฐ์ ํด์์ต๋๋ค. SharpSkill์ ์ด์ํ๋ฉฐ ์ด๊ณณ์ ๊ฒ์๋๋ ๋ชจ๋ ๋ด์ฉ์ ์ฑ ์์ ์ง๋๋ค.
2026๋ 4์ 20์ผ ์ ๋ฐ์ดํธ
ํ๊ทธ
๊ณต์
๊ด๋ จ ๊ธฐ์ฌ

2026๋ ๋ฐ์ดํฐ ์์ง๋์ด๋ง ๋ฉด์ ์ง๋ฌธ ์์ 25๊ฐ
2026๋ ๋ฐ์ดํฐ ์์ง๋์ด๋ง ๋ฉด์ ์์ ๊ฐ์ฅ ๋ง์ด ์ถ์ ๋๋ 25๊ฐ์ง ํต์ฌ ์ง๋ฌธ๊ณผ ์ค๋ฌด ์ค์ฌ์ ๋ต๋ณ์ ์ ๊ณตํฉ๋๋ค.

2026๋ Delta Lake vs Apache Iceberg: ๋ ์ดํฌํ์ฐ์ค ์ํคํ ์ฒ์ ๋ฉด์ ๋๋น ๊ฐ์ด๋
Delta Lake์ Apache Iceberg์ ๊ธฐ์ ์ ์ฐจ์ด์ ์ ์์ธํ ๋ถ์ํฉ๋๋ค. ํํฐ์ ์งํ, ACID ํธ๋์ญ์ , ์ฟผ๋ฆฌ ์์ง ํธํ์ฑ ๋ฑ ๋ฐ์ดํฐ ๋ ์ดํฌํ์ฐ์ค ๋ฉด์ ์์ ์์ฃผ ์ถ์ ๋๋ ์ฃผ์ ๋ฅผ ํฌ๊ด์ ์ผ๋ก ๋ค๋ฃน๋๋ค.

dbt 2026 ์๋ฒฝ ๊ฐ์ด๋: ๋ฐ์ดํฐ ๋ณํ, ํ ์คํธ ์ ๋ต, ๋ฉด์ ์ง๋ฌธ ์ด์ ๋ฆฌ
dbt๋ฅผ ํ์ฉํ ๋ฐ์ดํฐ ๋ณํ์ ํต์ฌ ๊ฐ๋ ๋ถํฐ ์ค๋ฌด๊น์ง, ๋ ์ด์ด๋ ๋ชจ๋ธ๋ง, ์ธํฌ๋ฆฌ๋ฉํ ์ ๋ต, ํ ์คํธ ๋ฐฉ๋ฒ๋ก , ๊ทธ๋ฆฌ๊ณ 2026๋ ๋ฐ์ดํฐ ์์ง๋์ด๋ง ๋ฉด์ ์์ ์์ฃผ ์ถ์ ๋๋ ์ง๋ฌธ์ ์ฝ๋ ์์ ์ ํจ๊ป ์์ธํ ๋ค๋ฃน๋๋ค.