๋ฐ์ดํ„ฐ ์—”์ง€๋‹ˆ์–ด๋ฅผ ์œ„ํ•œ Apache Kafka: ์ŠคํŠธ๋ฆฌ๋ฐ, ํŒŒํ‹ฐ์…˜, ๋ฉด์ ‘ ์งˆ๋ฌธ

๋ฐ์ดํ„ฐ ์—”์ง€๋‹ˆ์–ด๋ฅผ ์œ„ํ•œ Apache Kafka ์‹ฌ์ธต ๋ถ„์„. Kafka 4.x์™€ KRaft๋ฅผ ํ™œ์šฉํ•œ ์ŠคํŠธ๋ฆฌ๋ฐ ์•„ํ‚คํ…์ฒ˜, ํŒŒํ‹ฐ์…˜ ์ „๋žต, ์ปจ์Šˆ๋จธ ๊ทธ๋ฃน, ๊ธฐ์ˆ  ๋ฉด์ ‘ ๋นˆ์ถœ ์งˆ๋ฌธ์„ ์‹ค์ „ ์ฝ”๋“œ ์˜ˆ์ œ์™€ ํ•จ๊ป˜ ์„ค๋ช…ํ•ฉ๋‹ˆ๋‹ค.

Apache Kafka ์ŠคํŠธ๋ฆฌ๋ฐ ์•„ํ‚คํ…์ฒ˜์™€ ํŒŒํ‹ฐ์…˜ ๋ฐ์ดํ„ฐ ํ๋ฆ„ ๋‹ค์ด์–ด๊ทธ๋žจ

Apache Kafka๋Š” ํ˜„๋Œ€ ๋ฐ์ดํ„ฐ ์—”์ง€๋‹ˆ์–ด๋ง ์Šคํƒ์˜ ํ•ต์‹ฌ์— ์ž๋ฆฌํ•˜๋ฉฐ, ๋ชจ๋“  ๊ทœ๋ชจ์˜ ์กฐ์ง์—์„œ ๋งค์ผ ์ˆ˜์กฐ ๊ฑด์˜ ์ด๋ฒคํŠธ๋ฅผ ์ฒ˜๋ฆฌํ•˜๊ณ  ์žˆ์Šต๋‹ˆ๋‹ค. Kafka 4.x๋ถ€ํ„ฐ๋Š” KRaft ๊ธฐ๋ฐ˜์˜ ์™„์ „ ์ž์œจ ์šด์˜์ด ๊ตฌํ˜„๋˜์–ด ZooKeeper ์˜์กด์„ฑ์ด ์™„์ „ํžˆ ์ œ๊ฑฐ๋˜์—ˆ์Šต๋‹ˆ๋‹ค. ๋ฆฌ์–ผํƒ€์ž„ ๋ฐ์ดํ„ฐ ํŒŒ์ดํ”„๋ผ์ธ์˜ ์—…๊ณ„ ํ‘œ์ค€์œผ๋กœ์„œ ๋™์ผํ•œ ์•ˆ์ •์„ฑ์„ ์œ ์ง€ํ•˜๋ฉด์„œ๋„, ์šด์˜ ๋ณต์žก์„ฑ์ด ํฌ๊ฒŒ ๊ฐ์†Œํ–ˆ์Šต๋‹ˆ๋‹ค.

Kafka 4.x: KRaft ์ „์šฉ ๋ชจ๋“œ

Apache Kafka 4.0๋ถ€ํ„ฐ ZooKeeper๊ฐ€ ๋” ์ด์ƒ ํ•„์š”ํ•˜์ง€ ์•Š์Šต๋‹ˆ๋‹ค. KRaft(Kafka Raft)๊ฐ€ ๋ชจ๋“  ๋ฉ”ํƒ€๋ฐ์ดํ„ฐ ๊ด€๋ฆฌ๋ฅผ ๋„ค์ดํ‹ฐ๋ธŒ๋กœ ์ฒ˜๋ฆฌํ•˜์—ฌ ์šด์˜ ๋ณต์žก์„ฑ์„ ์ค„์ด๊ณ , ํ”„๋กœ๋•์…˜ ๋ฐฐํฌ์—์„œ ์ธํ”„๋ผ ๊ตฌ์„ฑ ์š”์†Œ๋ฅผ ํ•˜๋‚˜ ์™„์ „ํžˆ ์ œ๊ฑฐํ–ˆ์Šต๋‹ˆ๋‹ค.

๋ฐ์ดํ„ฐ ํŒŒ์ดํ”„๋ผ์ธ์„ ์œ„ํ•œ Kafka ์ŠคํŠธ๋ฆฌ๋ฐ ์•„ํ‚คํ…์ฒ˜

Kafka๋Š” ๋ถ„์‚ฐ ์ปค๋ฐ‹ ๋กœ๊ทธ๋กœ ๋™์ž‘ํ•ฉ๋‹ˆ๋‹ค. ํ”„๋กœ๋“€์„œ๊ฐ€ ํ† ํ”ฝ ํŒŒํ‹ฐ์…˜์˜ ๋์— ๋ ˆ์ฝ”๋“œ๋ฅผ ์ถ”๊ฐ€ํ•˜๊ณ , ์ปจ์Šˆ๋จธ๊ฐ€ ํ•ด๋‹น ๋ ˆ์ฝ”๋“œ๋ฅผ ์ˆœ์„œ๋Œ€๋กœ ์ฝ์Šต๋‹ˆ๋‹ค. ์ด ์ถ”๊ฐ€ ์ „์šฉ(append-only) ์„ค๊ณ„ ๋•๋ถ„์— ๋ฆฌ์–ผํƒ€์ž„ ์ŠคํŠธ๋ฆฌ๋ฐ๊ณผ ์ž„์˜ ์˜คํ”„์…‹์—์„œ์˜ ๋ฐฐ์น˜ ๋ฆฌํ”Œ๋ ˆ์ด๊ฐ€ ๋ชจ๋‘ ๊ฐ€๋Šฅํ•˜๋ฉฐ, ๋‘ ๊ฐ€์ง€ ํŒจํ„ด์„ ๋ชจ๋‘ ํ•„์š”๋กœ ํ•˜๋Š” ๋ฐ์ดํ„ฐ ์—”์ง€๋‹ˆ์–ด๋ง ์›Œํฌ๋กœ๋“œ์— ์ตœ์ ์˜ ๊ธฐ๋ฐ˜์„ ์ œ๊ณตํ•ฉ๋‹ˆ๋‹ค.

Kafka ํด๋Ÿฌ์Šคํ„ฐ๋Š” ๋ธŒ๋กœ์ปค(์„œ๋ฒ„), ํ† ํ”ฝ(๋…ผ๋ฆฌ์  ์ฑ„๋„), ํŒŒํ‹ฐ์…˜(ํ† ํ”ฝ์˜ ๋ฌผ๋ฆฌ์  ์ƒค๋“œ)์œผ๋กœ ๊ตฌ์„ฑ๋ฉ๋‹ˆ๋‹ค. ๊ฐ ํŒŒํ‹ฐ์…˜์€ ์ˆœ์„œ๊ฐ€ ์ง€์ •๋œ ๋ถˆ๋ณ€์˜ ๋ ˆ์ฝ”๋“œ ์‹œํ€€์Šค์ž…๋‹ˆ๋‹ค. ํ•˜๋‚˜์˜ ๋ธŒ๋กœ์ปค๊ฐ€ ํŒŒํ‹ฐ์…˜ ๋ฆฌ๋” ์—ญํ• ์„ ํ•˜๋ฉฐ ๋ชจ๋“  ์ฝ๊ธฐ์™€ ์“ฐ๊ธฐ๋ฅผ ์ฒ˜๋ฆฌํ•ฉ๋‹ˆ๋‹ค. ํŒ”๋กœ์›Œ ๋ ˆํ”Œ๋ฆฌ์นด๋Š” ๋‚ด๊ฒฐํ•จ์„ฑ์„ ์œ„ํ•ด ์‚ฌ๋ณธ์„ ์œ ์ง€ํ•ฉ๋‹ˆ๋‹ค.

ํ•ต์‹ฌ ์•„ํ‚คํ…์ฒ˜ ์†์„ฑ์€ ๋‹ค์Œ๊ณผ ๊ฐ™์Šต๋‹ˆ๋‹ค: ์ˆœ์„œ ๋ณด์žฅ์€ ํŒŒํ‹ฐ์…˜ ๋‚ด์—์„œ๋งŒ ์œ ํšจํ•˜๋ฉฐ, ํŒŒํ‹ฐ์…˜ ๊ฐ„์—๋Š” ๋ณด์žฅ๋˜์ง€ ์•Š์Šต๋‹ˆ๋‹ค. ์ด ๋‹จ์ผ ์ œ์•ฝ ์กฐ๊ฑด์ด ๋ฐ์ดํ„ฐ ํŒŒ์ดํ”„๋ผ์ธ์—์„œ ๋Œ€๋ถ€๋ถ„์˜ Kafka ์„ค๊ณ„ ๊ฒฐ์ •์„ ์ขŒ์šฐํ•ฉ๋‹ˆ๋‹ค.

yaml
# 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์˜ ๋ฉ”ํƒ€๋ฐ์ดํ„ฐ ๊ฐœ์„  ๋•๋ถ„์— ๋ธŒ๋กœ์ปค๋‹น ์ˆ˜์ฒœ ๊ฐœ์˜ ํŒŒํ‹ฐ์…˜์„ ํšจ์œจ์ ์œผ๋กœ ์ฒ˜๋ฆฌํ•ฉ๋‹ˆ๋‹ค

ํŒŒํ‹ฐ์…˜ ํ‚ค ์„ ํƒ์— ๋”ฐ๋ผ ์–ด๋–ค ๋ ˆ์ฝ”๋“œ๊ฐ€ ์–ด๋–ค ํŒŒํ‹ฐ์…˜์— ๋ฐฐ์น˜๋ ์ง€ ๊ฒฐ์ •๋ฉ๋‹ˆ๋‹ค. ๊ฐ™์€ ํ‚ค๋ฅผ ๊ฐ€์ง„ ๋ ˆ์ฝ”๋“œ๋Š” ํ•ญ์ƒ ๊ฐ™์€ ํŒŒํ‹ฐ์…˜์œผ๋กœ ์ „์†ก๋˜์–ด, ํ•ด๋‹น ํ‚ค์— ๋Œ€ํ•œ ์ˆœ์„œ๊ฐ€ ๋ณด์กด๋ฉ๋‹ˆ๋‹ค.

python
# 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 ํ† ํ”ฝ์— ์ €์žฅํ•ฉ๋‹ˆ๋‹ค. ์ฒ˜๋ฆฌ ๋Œ€๋น„ ์˜คํ”„์…‹ ์ปค๋ฐ‹์˜ ํƒ€์ด๋ฐ์ด ์ „๋‹ฌ ๋ณด์žฅ ์ˆ˜์ค€์„ ๊ฒฐ์ •ํ•ฉ๋‹ˆ๋‹ค.

python
# 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.)
    pass

enable.auto.commit๋ฅผ False๋กœ ์„ค์ •ํ•˜๊ณ  write_to_warehouse ์„ฑ๊ณต ํ›„์—๋งŒ ์ปค๋ฐ‹ํ•จ์œผ๋กœ์จ At-least-once ์ „๋‹ฌ์ด ๋ณด์žฅ๋ฉ๋‹ˆ๋‹ค. ์ปจ์Šˆ๋จธ๊ฐ€ ์ฒ˜๋ฆฌ ํ›„ ์ปค๋ฐ‹ ์ „์— ํฌ๋ž˜์‹œ๋˜๋ฉด, ์žฌ์‹œ์ž‘ ์‹œ ๋ ˆ์ฝ”๋“œ๊ฐ€ ์žฌ์ „๋‹ฌ๋ฉ๋‹ˆ๋‹ค. ์›จ์–ดํ•˜์šฐ์Šค์— ๋ฐ์ดํ„ฐ๋ฅผ ๊ณต๊ธ‰ํ•˜๋Š” ํŒŒ์ดํ”„๋ผ์ธ์—์„œ๋Š” ์ด๊ฒƒ์ด ์ผ๋ฐ˜์ ์œผ๋กœ ์˜ฌ๋ฐ”๋ฅธ ํŠธ๋ ˆ์ด๋“œ์˜คํ”„์ž…๋‹ˆ๋‹ค. ์›จ์–ดํ•˜์šฐ์Šค์— ๋Œ€ํ•œ ๋ฉฑ๋“ฑ ์“ฐ๊ธฐ๊ฐ€ ์ค‘๋ณต์„ ์ฒ˜๋ฆฌํ•ฉ๋‹ˆ๋‹ค.

๋ฐ์ดํ„ฐ ํŒŒ์ดํ”„๋ผ์ธ ํ†ตํ•ฉ์„ ์œ„ํ•œ Kafka Connect

Kafka Connect๋Š” ์ปค์Šคํ…€ ํ”„๋กœ๋“€์„œ/์ปจ์Šˆ๋จธ ์ฝ”๋“œ๋ฅผ ์ž‘์„ฑํ•˜์ง€ ์•Š๊ณ ๋„ Kafka์™€ ์™ธ๋ถ€ ์‹œ์Šคํ…œ ๊ฐ„์˜ ๋ฐ์ดํ„ฐ ์ด๋™์„ ๊ฐ€๋Šฅํ•˜๊ฒŒ ํ•˜๋Š” ํ™•์žฅ ๊ฐ€๋Šฅํ•œ ํ”„๋ ˆ์ž„์›Œํฌ์ž…๋‹ˆ๋‹ค. ์†Œ์Šค ์ปค๋„ฅํ„ฐ๊ฐ€ Kafka๋กœ ๋ฐ์ดํ„ฐ๋ฅผ ์ˆ˜์ง‘ํ•˜๊ณ , ์‹ฑํฌ ์ปค๋„ฅํ„ฐ๊ฐ€ ๋ฐ์ดํ„ฐ๋ฅผ ์™ธ๋ถ€๋กœ ์ „์†กํ•ฉ๋‹ˆ๋‹ค.

json
{
  "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๊ฐ€์ง€ ๊ตฌ์„ฑ ์š”์†Œ:

  1. ๋ฉฑ๋“ฑ ํ”„๋กœ๋“€์„œ: ๊ฐ ํ”„๋กœ๋“€์„œ๋Š” ๊ณ ์œ ํ•œ Producer ID๋ฅผ ๋ถ€์—ฌ๋ฐ›์Šต๋‹ˆ๋‹ค. ๋ธŒ๋กœ์ปค๋Š” ์‹œํ€€์Šค ๋ฒˆํ˜ธ๋ฅผ ์‚ฌ์šฉํ•˜์—ฌ ์žฌ์‹œ๋„์˜ ์ค‘๋ณต์„ ์ œ๊ฑฐํ•˜๋ฏ€๋กœ, ๋„คํŠธ์›Œํฌ ์žฌ์‹œ๋„๊ฐ€ ์ค‘๋ณต์„ ์ƒ์„ฑํ•˜์ง€ ์•Š์Šต๋‹ˆ๋‹ค.
  2. ํŠธ๋žœ์žญ์…˜: ํ”„๋กœ๋“€์„œ๋Š” ์—ฌ๋Ÿฌ ํŒŒํ‹ฐ์…˜์— ๋Œ€ํ•œ ์“ฐ๊ธฐ์™€ ์˜คํ”„์…‹ ์ปค๋ฐ‹์„ ๋‹จ์ผ ํŠธ๋žœ์žญ์…˜์œผ๋กœ ์›์ž์ ์œผ๋กœ ์ˆ˜ํ–‰ํ•  ์ˆ˜ ์žˆ์Šต๋‹ˆ๋‹ค. ๋ชจ๋“  ์“ฐ๊ธฐ๊ฐ€ ์„ฑ๊ณตํ•˜๊ฑฐ๋‚˜ ๋ชจ๋‘ ์‹คํŒจํ•ฉ๋‹ˆ๋‹ค.
  3. read_committed ๊ฒฉ๋ฆฌ ์ˆ˜์ค€: isolation.level=read_committed๋กœ ์„ค์ •๋œ ์ปจ์Šˆ๋จธ๋Š” ์ปค๋ฐ‹๋œ ํŠธ๋žœ์žญ์…˜์˜ ๋ ˆ์ฝ”๋“œ๋งŒ ์ฝ์Šต๋‹ˆ๋‹ค.
python
# 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 record

transactional.id ์„ค์ •์„ ํ†ตํ•ด ๋ธŒ๋กœ์ปค๋Š” ์ข€๋น„ ํ”„๋กœ๋“€์„œ๋ฅผ ์ฐจ๋‹จํ•  ์ˆ˜ ์žˆ์Šต๋‹ˆ๋‹ค. ์ปจ์Šˆ๋จธ๊ฐ€ ํฌ๋ž˜์‹œ๋˜๊ณ  ๋™์ผํ•œ transactional.id๋กœ ์ƒˆ ์ธ์Šคํ„ด์Šค๊ฐ€ ์‹œ์ž‘๋˜๋ฉด, ๋ธŒ๋กœ์ปค๋Š” ์ด์ „ ์ธ์Šคํ„ด์Šค์˜ ๋ณด๋ฅ˜ ์ค‘์ธ ํŠธ๋žœ์žญ์…˜์„ ์ค‘๋‹จํ•˜์—ฌ ์ถœ๋ ฅ ์ค‘๋ณต์„ ๋ฐฉ์ง€ํ•ฉ๋‹ˆ๋‹ค.

Share Groups: Kafka 4.2์˜ ํ ์‹œ๋งจํ‹ฑ

Kafka 4.2์—์„œ๋Š” ํ”„๋กœ๋•์…˜ ์ˆ˜์ค€์˜ Share Groups๊ฐ€ ๋„์ž…๋˜์–ด, Kafka์— ์ „ํ†ต์ ์ธ ๋ฉ”์‹œ์ง€ ํ ์‹œ๋งจํ‹ฑ์ด ์ถ”๊ฐ€๋˜์—ˆ์Šต๋‹ˆ๋‹ค. ๊ฐ ํŒŒํ‹ฐ์…˜์ด ๊ทธ๋ฃน ๋‚ด ์ •ํ™•ํžˆ ํ•˜๋‚˜์˜ ์ปจ์Šˆ๋จธ์—๋งŒ ํ• ๋‹น๋˜๋Š” ์ปจ์Šˆ๋จธ ๊ทธ๋ฃน๊ณผ ๋‹ฌ๋ฆฌ, Share Groups์—์„œ๋Š” ์—ฌ๋Ÿฌ ์ปจ์Šˆ๋จธ๊ฐ€ ๋™์ผํ•œ ํŒŒํ‹ฐ์…˜์˜ ๋ ˆ์ฝ”๋“œ๋ฅผ ๋…๋ฆฝ์ ์œผ๋กœ ์ฒ˜๋ฆฌํ•  ์ˆ˜ ์žˆ์Šต๋‹ˆ๋‹ค.

ํŠน์„ฑ์ปจ์Šˆ๋จธ ๊ทธ๋ฃนShare Groups
ํŒŒํ‹ฐ์…˜ ํ• ๋‹น๋ฐฐํƒ€์  (ํŒŒํ‹ฐ์…˜๋‹น ์ปจ์Šˆ๋จธ 1๊ฐœ)๊ณต์œ  (ํŒŒํ‹ฐ์…˜๋‹น ๋ณต์ˆ˜ ์ปจ์Šˆ๋จธ)
์ˆœ์„œ ๋ณด์žฅํŒŒํ‹ฐ์…˜ ๋‚ด ์ˆœ์„œ ๋ณด์กด์ˆœ์„œ ๋ณด์žฅ ์—†์Œ
์‚ฌ์šฉ ์‚ฌ๋ก€์ˆœ์„œ ๋ณด์žฅ ์ด๋ฒคํŠธ ์ŠคํŠธ๋ฆผ์ž‘์—… ๋ถ„๋ฐฐ, ์ž‘์—… ํ
ํ™•์ธ ์‘๋‹ต์˜คํ”„์…‹ ๊ธฐ๋ฐ˜ (์œ„์น˜ ์ปค๋ฐ‹)๋ ˆ์ฝ”๋“œ ๋‹จ์œ„ (๊ฐœ๋ณ„ ํ™•์ธ ์‘๋‹ต)
์ตœ๋Œ€ ๋ณ‘๋ ฌ ์ฒ˜๋ฆฌํŒŒํ‹ฐ์…˜ ์ˆ˜๋กœ ์ œํ•œ์ปจ์Šˆ๋จธ ์ˆ˜๋กœ ์ œํ•œ

Share Groups๋Š” ์˜ค๋žซ๋™์•ˆ ์กด์žฌํ•˜๋˜ ํ•œ๊ณ„, ์ฆ‰ ํŒŒํ‹ฐ์…˜ ์ˆ˜๋ฅผ ์ดˆ๊ณผํ•˜๋Š” ์ปจ์Šˆ๋จธ ํ™•์žฅ ๋ฌธ์ œ๋ฅผ ํ•ด๊ฒฐํ•ฉ๋‹ˆ๋‹ค. ์•Œ๋ฆผ ์ „์†ก์ด๋‚˜ ๋ฐฐ์น˜ ์ž‘์—… ๋ถ„๋ฐฐ์™€ ๊ฐ™์ด ์ˆœ์„œ๊ฐ€ ํ•„์š” ์—†๋Š” ์›Œํฌ๋กœ๋“œ์—์„œ๋Š” Kafka ์™ธ์— ๋ณ„๋„์˜ ์™ธ๋ถ€ ํ ์‹œ์Šคํ…œ์ด ํ•„์š”ํ•˜์ง€ ์•Š๊ฒŒ ๋ฉ๋‹ˆ๋‹ค.

Share Groups ํ™œ์šฉ ์‹œ์ 

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-Maillet

์ž‘์„ฑ์ž

Anthony Fillion-Maillet

SharpSkill ์ฐฝ์—…์ž

10๋…„ ์ด์ƒ ํ’€์Šคํƒ ๊ฐœ๋ฐœ์„ ํ•ด์™”์Šต๋‹ˆ๋‹ค. SharpSkill์„ ์šด์˜ํ•˜๋ฉฐ ์ด๊ณณ์— ๊ฒŒ์‹œ๋˜๋Š” ๋ชจ๋“  ๋‚ด์šฉ์— ์ฑ…์ž„์„ ์ง‘๋‹ˆ๋‹ค.

2026๋…„ 4์›” 20์ผ ์—…๋ฐ์ดํŠธ

ํƒœ๊ทธ

#kafka
#streaming
#data-engineering
#partitions
#consumer-groups
#kraft
#event-driven

๊ณต์œ 

๊ด€๋ จ ๊ธฐ์‚ฌ

Data Engineering Interview Questions 2026

2026๋…„ ๋ฐ์ดํ„ฐ ์—”์ง€๋‹ˆ์–ด๋ง ๋ฉด์ ‘ ์งˆ๋ฌธ ์ƒ์œ„ 25๊ฐœ

2026๋…„ ๋ฐ์ดํ„ฐ ์—”์ง€๋‹ˆ์–ด๋ง ๋ฉด์ ‘์—์„œ ๊ฐ€์žฅ ๋งŽ์ด ์ถœ์ œ๋˜๋Š” 25๊ฐ€์ง€ ํ•ต์‹ฌ ์งˆ๋ฌธ๊ณผ ์‹ค๋ฌด ์ค‘์‹ฌ์˜ ๋‹ต๋ณ€์„ ์ œ๊ณตํ•ฉ๋‹ˆ๋‹ค.

2026๋…„ Delta Lake vs Apache Iceberg: ๋ ˆ์ดํฌํ•˜์šฐ์Šค ์•„ํ‚คํ…์ฒ˜์™€ ๋ฉด์ ‘ ๋Œ€๋น„ ๊ฐ€์ด๋“œ

2026๋…„ Delta Lake vs Apache Iceberg: ๋ ˆ์ดํฌํ•˜์šฐ์Šค ์•„ํ‚คํ…์ฒ˜์™€ ๋ฉด์ ‘ ๋Œ€๋น„ ๊ฐ€์ด๋“œ

Delta Lake์™€ Apache Iceberg์˜ ๊ธฐ์ˆ ์  ์ฐจ์ด์ ์„ ์ƒ์„ธํžˆ ๋ถ„์„ํ•ฉ๋‹ˆ๋‹ค. ํŒŒํ‹ฐ์…˜ ์ง„ํ™”, ACID ํŠธ๋žœ์žญ์…˜, ์ฟผ๋ฆฌ ์—”์ง„ ํ˜ธํ™˜์„ฑ ๋“ฑ ๋ฐ์ดํ„ฐ ๋ ˆ์ดํฌํ•˜์šฐ์Šค ๋ฉด์ ‘์—์„œ ์ž์ฃผ ์ถœ์ œ๋˜๋Š” ์ฃผ์ œ๋ฅผ ํฌ๊ด„์ ์œผ๋กœ ๋‹ค๋ฃน๋‹ˆ๋‹ค.

dbt data transformations and testing tutorial 2026

dbt 2026 ์™„๋ฒฝ ๊ฐ€์ด๋“œ: ๋ฐ์ดํ„ฐ ๋ณ€ํ™˜, ํ…Œ์ŠคํŠธ ์ „๋žต, ๋ฉด์ ‘ ์งˆ๋ฌธ ์ด์ •๋ฆฌ

dbt๋ฅผ ํ™œ์šฉํ•œ ๋ฐ์ดํ„ฐ ๋ณ€ํ™˜์˜ ํ•ต์‹ฌ ๊ฐœ๋…๋ถ€ํ„ฐ ์‹ค๋ฌด๊นŒ์ง€, ๋ ˆ์ด์–ด๋“œ ๋ชจ๋ธ๋ง, ์ธํฌ๋ฆฌ๋ฉ˜ํƒˆ ์ „๋žต, ํ…Œ์ŠคํŠธ ๋ฐฉ๋ฒ•๋ก , ๊ทธ๋ฆฌ๊ณ  2026๋…„ ๋ฐ์ดํ„ฐ ์—”์ง€๋‹ˆ์–ด๋ง ๋ฉด์ ‘์—์„œ ์ž์ฃผ ์ถœ์ œ๋˜๋Š” ์งˆ๋ฌธ์„ ์ฝ”๋“œ ์˜ˆ์ œ์™€ ํ•จ๊ป˜ ์ƒ์„ธํžˆ ๋‹ค๋ฃน๋‹ˆ๋‹ค.