Kafka Consumer (1)

Ganplank·2025년 10월 9일
post-thumbnail

1. Kafka Consumer?

  • Kafka Topic 에 있는 메시지를 읽어서 소비해주는 역할
    Offset내용을 commit 하여 이전까지의 메시지는 처리 완료 상태를 Kafka Broker에 기록한다.
    이를 통해 Consumer는 재시작 후에도 마지막 커밋 지점 이후부터 메시지를 이어서 소비할 수 있다.
  • Consumer는 동일한 group.id 를 가진 인스턴스들이 하나의 Consumer Group 으로 논리적으로 묶인다. 각 Consumer는 bootstrap.servers (doa-kafka-kafka-bootstrap:9092)를 통해 Kafka 클러스터에 접속하고,
    해당 Topic의 메시지를 그룹 단위로 분담하여 소비한다.
    • Partition보다 Consumer가 많아지면 남은 Consumer가 Partition 을 할당받지 못함
  • Kafka Broker의 Topic의 Partition들은 Consumer group의 각 인스턴스들과 1:1 매핑되고 각 Partition 내 Record 전송 순서는 유지되나 2개이상의 Parition 경우 Publisher가 순서 보장없이 병렬처리하여 순서가 보장되지 않는다.
    순서보장이 필요한 주문시스템 같은 경우 별도 Key를 정의해 해시한 값을 기준으로 특정 Partition 으로만 Record 전송하도록 하는 등 조치가 필요하다.
  • Publisher가 순서보장없이 Partition으로 메세지를 병렬 전송처리한다는 의미는 데이터 자체를 Segement해서 병렬전송하는 아니고 큰 한덩어리(Kafka Broker TOPIC BATCH SIZE에 담긴 레코드)를 순서보장없이 Partition에 적재한다는 의미로 이해된다.
  • Partition 개수는 분배방식이나 개수는 전략적으로 선정해야하며 Partition 개수 증가는 마음대로 가능하지만 줄이는것은 Topic 을 삭제하고 다시 생성해야하므로 Partition 증가할 때 여러가지 부분들을 고려해야 한다.

2. RMA 자동화 Consumer 역할은?

1) 설명

  • Consumer는 upserter.py 로 정의되어있다.
  • Upserter Consumer 는 Consumer 역할 뿐만 아니라 Producer 역할도 하는데 Commit 실패 등 Google API 전송 실패 시 오류 내용을 DLQ Topic에 Message 를 전달하는 역할도 한다.
  • Topic(cdc.public.rma_lists)에 Message는 Postgresql PUBLICATION 등록되어있는 rma_lists 테이블 업데이트 시 Replication Slot을 통해 Debezium CDC가 WAL을 읽어 Kafka Topic에 Message를 전달한다. 해당 Message를 Consumer가 읽어서 Goolge API 전송 후 정상처리 완료시에만 Kafka Broker에 Offset을 늘리고 Commit한다.

2) 설정

  • Kafka-Topic-rma.lists

    apiVersion: kafka.strimzi.io/v1beta2
     kind: KafkaTopic
     metadata:
       name: cdc.public.rma-lists
       namespace: streaming
       labels:
         strimzi.io/cluster: doa-kafka   # 클러스터 이름 Label
     spec:
       topicName: cdc.public.rma_lists
       partitions: 3                     # 병렬 소비 Parition 개수
       replicas: 1                       # Kafka Broker 보다 클 수 없음 
         config:				    		# Topic Message보관관련
         retention.ms: 604800000         ## 7일
         segment.bytes: 1073741824       ## 1GB
  • Consumer.group.id

    def build_consumer()->Consumer:
     return Consumer({
         "bootstrap.servers": BOOTSTRAP, ## doa-kafka-kafka-bootstrap:9092
         "group.id": KAFKA_GROUP_ID, ## sheets-upserter
    			....
     })
    • spec.relicas: 3 으로 partition 개수 맞춤
      apiVersion: apps/v1
      kind: Deployment
      metadata:
        name: sheets-upserter
        namespace: streaming
      spec:
        replicas: 3
        selector:
          matchLabels: { app: sheets-upserter }
        template:
          metadata:
            labels: { app: sheets-upserter }
            annotations:
              prometheus.io/scrape: "true"
              prometheus.io/port: "8000"
              prometheus.io/path: "/metrics"
          spec:
            containers:
            - name: app
              image: ganplank/sheets-upserter:10.2
  • TOPIC 내용 확인

    • kconsume 명령어로 cdc.public.rma_lists TOPIC 이 들어왔는지 로그를 통해 확인 가능
    • kafka brokers 가 제공해주는 consumer.sh 을 통해 확인 지원
    alias kconsume='kubectl -n streaming exec -it doa-kafka-kp-brokers-0 -- /opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic cdc.public.rma_lists --from-beginning'
    • kconsume 명령어 실행 시 로그 내용
      • Ticket ID 선정 기준 : sha256("rma|ts|" + random(16B) + "|" + pepper)
    {"id":"rma-20251013061548-SwzL3S","name":"흠냐","company":"흠냐","email":"hmm@gmail.com","model":"kaka","serial_number":"ka","as_method":"유상","initial_install_date":20367,"failure_date":null,"version":null,"memo":"kaka","status":"pending","created_at":1760336148304617,"__deleted":"false","__op":"c","__source_ts_ms":1760336148305,"__table":"rma_lists"}
  • Consumer 수동 + 비동기 동작 설정

    def build_consumer()->Consumer:
      return Consumer({
          "bootstrap.servers": BOOTSTRAP,
          "group.id": KAFKA_GROUP_ID,
          "enable.auto.commit": False,
    
      ... 
    
      try:
        with m_batch_latency.time():
          apply_batch(state, upserts, deletes)
        consumer.commit(offsets=last_offsets, asynchronous=False)
        m_msgs_ok.inc(len(upserts)+len(deletes))
    • Google Sheets 업데이트 실패 시 Commit되지 않도록하여 장애상황에 메시지 분실 위험을 제거
    • 업데이트 실패 후 Message 재시도로 인해 중복 데이터가 생길 수있는 부분은 Google Sheets의 A Columne에 __PK 값 기준으로 멱등성을 보장하여 최신 Message 기준으로 업데이트되므로 중복제거
    • 행 수정 및 제거 또한 __PK 값 기준으로 수행됨
  • Kafka Consumer Loop 구조 에러 발생 시 재시도 설정

    @retry(wait=wait_exponential(min=1, max=30), stop=stop_after_attempt(6),
        retry=retry_if_exception_type((HttpError, ConnectionError, TimeoutError, BrokenPipeError)))
    • HttpError, ConnectionError, TimeoutError, BrokenPipeError 에러 경우
  • consumer Poll 설정

    BATCH_SIZE   = int(os.environ.get("BATCH_SIZE", "200"))
     POLL_MS      = int(os.environ.get("POLL_MS", "1000"))
     ....
        try:
          while not stop_flag:
              msgs = consumer.consume(BATCH_SIZE, timeout=POLL_MS/1000.0)
              if not msgs: continue
    
    • BATCH_SIZE : Message(Google Sheet 행) 단위
    • 응답대기시간 : POLL_MS/1000.0 = 초 단위
  • Prometheus 매트릭 수집 (진행 예정)

  • 이메일 및 슬랙 메시지 알림 (진행 예정)

3) 동작 순서

1) poll: 카프카에서 최대 BATCH_SIZE개 메시지 가져옴
2) apply: 모은 레코드를 한 번에 Sheets API(append/update/batchUpdate)로 반영
3) commit: 위 2번이 정상 완료됐을 때만, 그 배치의 마지막 오프셋까지 커밋

다음 배치로 go

  • consumer 동작 플로우

    ┌────────────────────────────┐
    │  PostgreSQL (Debezium)   │
    │  → cdc.public.rma_lists  │
    └──────────────┬─────────────┘
                  │  Kafka Broker
                  ▼
    ┌───────────────────────────────┐
    │  Kafka Consumer (upserter) │
    │  poll()                    │
    │  ↓                         │
    │  메시지 수신                │
    │  ↓                         │
    │  Google Sheets API 업데이트 │
    │  ↓                         │
    │  처리 성공 시:              │
    │  last_offsets.append(      │
    │     TopicPartition(        │
    │       msg.topic(),         │
    │       msg.partition(),     │
    │       msg.offset() + 1))   │
    │  ↓                         │
    │  consumer.commit(offsets=last_offsets) │
    │  ↓                         │
    │  Prometheus metrics inc()  │
    │  ↓                         │
    │  다음 poll() 실행           │
    └──────────────────────────────┘
                   │
              (예외 발생 시)
                   ▼
    ┌────────────────────────────┐
    │  Kafka DLQ Producer      │
    │  produce(DLQ_TOPIC, ...) │
    └────────────────────────────┘
  • 흐름표

    단계동작세부 내용
    ① Kafka Pollconsumer.poll()Kafka로부터 CDC 메시지 수신
    ② 파싱json.loads(msg.value())Debezium JSON이벤트를 Python Dictionary로 파싱
    ③ Sheets 업데이트update_sheets()Google Sheets API로 데이터 쓰기 (.values().update())
    ④ 성공 시 offset 기록msg.offset() + 1다음 메시지부터 읽도록 마킹
    ⑤ 커밋consumer.commit(offsets=last_offsets)Kafka Broker에게 처리완료 Commit되면 내부적으로 Kafka Broker __consumer_offsets Topic에 저장
    ⑥ 실패 시 DLQproducer.produce(DLQ_TOPIC, value=msg.value())오류난 메시지는 DLQ로 전송
    ⑦ 지표 증가m_api_calls.inc()Prometheus 메트릭 증가

3. Reference

https://sungjk.github.io/2021/01/10/kafka-consumer.html
https://pathtosenior.substack.com/p/a-gentle-introduction-to-kafka-consumer
https://velog.io/@jaehyeong/Apache-Kafka%EC%95%84%ED%8C%8C%EC%B9%98-%EC%B9%B4%ED%94%84%EC%B9%B4%EB%9E%80-%EB%AC%B4%EC%97%87%EC%9D%B8%EA%B0%80
https://velog.io/@hyun6ik/Apache-Kafka-Partition-Assignment-Strategy
https://yeongchan1228.tistory.com/m/60
https://curiousjinan.tistory.com/entry/understand-kafka-partitions
https://github.com/schooldevops/kafka-tutorials-with-kido/blob/main/01.kafka_install.md
https://www.jaenung.net/tree/28776

profile
안녕?

0개의 댓글