포스트

Kafka 기초 (5) - Python Client: confluent-kafka

confluent-kafka로 JSON 이벤트를 보내는 producer와 poll 루프로 읽는 consumer를 작성합니다. flush와 delivery callback, auto commit과 at-least-once까지 코드로 확인합니다.

Kafka 기초 (5) - Python Client: confluent-kafka

Kafka 기초 시리즈의 5편입니다. 2편의 kafka container가 떠 있는 상태에서 진행합니다. 전체 목차는 0편에 있습니다.

준비

console 도구로 하던 일을 애플리케이션 코드로 옮깁니다. Python client는 confluent-kafka 패키지를 씁니다. C 라이브러리(librdkafka) 기반이라 성능이 좋고 사실상 표준입니다.

1
2
3
mkdir kafka-py && cd kafka-py
python -m venv .venv && source .venv/bin/activate   # 가상환경은 python-env 글 참고
pip install confluent-kafka

폴더 구조는 파일 두 개가 전부입니다.

1
2
3
kafka-py/
├── producer.py
└── consumer.py

Producer

주문 이벤트를 JSON으로 보내는 producer입니다.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
import json
from confluent_kafka import Producer

producer = Producer({"bootstrap.servers": "localhost:9092"})

def on_delivery(err, msg):
    if err is not None:
        print(f"전송 실패: {err}")
    else:
        print(f"전송 완료: partition {msg.partition()}, offset {msg.offset()}")

for i in range(5):
    event = {"user_id": f"user-{i % 2}", "action": "click", "seq": i}
    producer.produce(
        "events",
        key=event["user_id"],          # 순서가 필요한 단위 = key (3편)
        value=json.dumps(event),
        on_delivery=on_delivery,
    )

producer.flush()   # 버퍼의 이벤트를 전부 전송하고 완료 확인까지 대기
1
python producer.py
1
2
3
4
전송 완료: partition 0, offset 12
전송 완료: partition 0, offset 13
전송 완료: partition 2, offset 8
...

설정과 흐름이 console 도구와 1:1로 대응합니다. bootstrap.servers--bootstrap-server이고, produce()의 topic·key·value가 producer 터미널에 치던 key:value 입력입니다. 코드에서 새로 등장하는 것은 두 가지입니다.

  • produce()는 보내지 않습니다. 이벤트를 client 내부 버퍼에 넣을 뿐이고, 전송은 백그라운드에서 모아서(batch) 일어납니다. 그래서 마지막에 flush()로 버퍼를 비우고 완료를 기다려야 합니다. 이것을 빼먹으면 프로그램이 이벤트를 보내기 전에 끝나 버리는, Python client의 가장 흔한 실수가 됩니다
  • 전송 결과는 callback으로 받습니다. on_delivery는 이벤트별 전송 성공·실패 시점에 호출됩니다. 출력에서 user-0과 user-1이 각자 한 partition에 몰리는 것도 확인됩니다 — 3편의 key 규칙이 코드에서도 그대로입니다

Consumer

poll 루프로 이벤트를 읽는 consumer입니다.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
import json
from confluent_kafka import Consumer

consumer = Consumer({
    "bootstrap.servers": "localhost:9092",
    "group.id": "loader",              # 4편의 --group
    "auto.offset.reset": "earliest",   # 이 group의 저장된 offset이 없을 때 처음부터
})
consumer.subscribe(["events"])

try:
    while True:
        msg = consumer.poll(1.0)       # 최대 1초 기다려 이벤트 하나를 반환
        if msg is None:                # 그 사이 새 이벤트가 없었음
            continue
        if msg.error():
            print(f"에러: {msg.error()}")
            continue
        event = json.loads(msg.value())
        print(f"partition {msg.partition()}, offset {msg.offset()}: {event}")
finally:
    consumer.close()   # group에서 정상 탈퇴 → rebalance가 즉시 일어난다
1
python consumer.py
1
2
3
4
partition 0, offset 12: {'user_id': 'user-0', 'action': 'click', 'seq': 0}
partition 0, offset 13: {'user_id': 'user-0', 'action': 'click', 'seq': 2}
partition 2, offset 8: {'user_id': 'user-1', 'action': 'click', 'seq': 1}
...
  • poll(1.0)이 심장입니다. consumer는 broker가 밀어주는 것을 받는 구조가 아니라, 루프를 돌며 직접 가져오는(poll) 구조입니다. 이 루프가 살아 있다는 것 자체가 group에 보내는 생존 신호를 겸해서, 루프가 오래 멈추면 죽은 것으로 간주되어 partition을 뺏깁니다
  • auto.offset.reset은 2편의 –from-beginning과 같은 문제의 답입니다. 이 group이 처음이라 저장된 offset이 없을 때 어디서 시작할지이며, earliest가 처음부터, 기본값 latest가 접속 이후부터입니다. 한 번 offset이 저장된 뒤에는 이 설정과 무관하게 저장 지점부터 이어 읽습니다
  • consumer를 하나 더 실행하면 4편의 partition 분담이 코드에서도 그대로 일어납니다


Offset Commit과 at-least-once

읽은 위치는 언제 저장되는 것일까요. 기본 설정(enable.auto.commit=true)에서는 client가 5초 간격으로 “여기까지 읽었다”를 자동 커밋합니다.

편하지만 정확히 한 번 처리를 보장하지는 않습니다. 이벤트를 처리한 직후, 커밋 전에 프로세스가 죽으면 재시작 후 그 이벤트를 다시 받습니다. 즉 기본 동작은 at-least-once(최소 한 번, 중복 가능)입니다. 반대로 “받자마자 커밋되고 처리 전에 죽는” 유실 방향은, 커밋 시점을 처리 완료 뒤로 직접 옮겨 막습니다.

1
2
3
4
5
6
7
consumer = Consumer({
    ...,
    "enable.auto.commit": False,
})
...
        # 처리(예: DB 적재)가 끝난 뒤에 커밋
        consumer.commit(msg)

그래도 “처리 완료 후, 커밋 전 사망” 구간의 중복은 남습니다. 그래서 실무의 원칙은 커밋 타이밍을 조이는 것에 더해 소비 쪽을 멱등하게 만드는 것입니다. 같은 이벤트가 두 번 와도 결과가 같도록 PK 기준 upsert로 적재하는 식이며, CDC 파이프라인에서 sink가 upsert 모드인 것이 정확히 이 이유입니다.

정리

console 도구 (2~4편)Python client
--bootstrap-serverbootstrap.servers
producer 터미널의 key:value 입력produce(topic, key=, value=) + flush()
--group loadergroup.id
--from-beginningauto.offset.reset=earliest
Ctrl+Cconsumer.close()

다음 편이 마지막입니다. broker가 죽어도 유실이 없으려면 무엇이 필요한지, retention과 advertised listener 같은 운영 설정과 자주 밟는 함정을 정리합니다.

다음 글: Kafka 기초 (6) - Operations Basics: 유실 방지, Retention, 함정

이 기사는 저작권자의 CC BY 4.0 라이센스를 따릅니다.