Kafka Consumer (2)

Ganplank·2025년 10월 22일
post-thumbnail

1. Notifier Consumer 역할?

  • Notifier Consumer 는 RMA DB업데이트 발생 시 관련 내용을 이메일과 슬랙 메시지 전송으로 관련 담당자에게 Notify해주는 Consumer이다.
  • 이전 구글스프레드시트 업데이트 Consumer와 유사하게 RMA 내용을 notify.rma KafkaTopic으로 메시지를 전송하면 Notifier Consumer가 Topic(Partition)에 메시지를 읽고 메일과 슬랙 메시지 전송 작업을한다.

2. Notifier 프로세스 추가

1) RMA 신청 WEB 생성

2) K8S SVC 생성

  • minikube환경에서 NodePort Service 생성 관련 YAML

3) 포트 포워딩 설정

  • RMA 웹페이지 접속을 위해서 구성이 현재 minikube 환경이므로 service IP가 NodeIP(물리서버)로 설정되지 않아 포트포워딩을 통해서 서버로 유입되는 트래픽을 minikube로 포워딩한다.
  • TCP/port 80은 Privileged Port이기 때문에 root 권한으로 포트포워딩 실행
    nohup sudo kubectl -n rma port-forward --address 0.0.0.0 svc/rma-web-svc 80:80 --kubeconfig $HOME/.kube/config &
    • 추후 Ingress-controll + minikbe tunnel 변경 예정이고 만약 물리서버가 증원된다면 K8S 클러스터링(Master 3, Worker 3)환경에서 배포 예정

4) Kafka Topic 생성

kind: KafkaTopic
metadata:
  name: notify.rma
  namespace: streaming
  labels:
    strimzi.io/cluster: doa-kafka
spec:
  partitions: 3
  replicas: 1
  config:
    cleanup.policy: "delete"
    retention.ms: 604800000        # 7 days
    min.insync.replicas: "1"
    segment.bytes: 1073741824      # 1 GiB

5) Produce 코드추가

  • uperter.py는 앞서 구글 스프레드시트 업데이트를 위해 읽어드린 Topic(cdc.public.customers)내용을 notify를 위한 Topic(notify.rma)으로 메세지를 보내줘야 한다.

    • upserter.py파일에 Produce해주는 코드 추가

      producer = Producer({"bootstrap.servers": BOOTSTRAP}) if (DLQ_TOPIC or NOTIFY_TOPIC) else None
      ...
      def send_notify(email: str, ticket_id: str, record: dict):
      
      tid = (ticket_id or record.get("id") or "").strip()
      to_email = (email or record.get("email") or "").strip()
      
      payload = {
          "name": record.get("name"),
          "model": record.get("model"),
          "serial_number": record.get("serial_number"),
          "email": to_email,
          "ticket_id": tid,
          "status": record.get("status"),
          "submitted_at_ms": record.get("__source_ts_ms"),
      }
      
      while True:  # 전송 에러 시 재시도 루프
          try:
              producer.produce(  ##NOTIFY_TOPIC Partition key 설정
                  NOTIFY_TOPIC,
                  key=tid.encode("utf-8"),
                  value=json.dumps(payload, ensure_ascii=False).encode("utf-8"),
                  on_delivery=_on_delivery            ## callback
              )
      
              producer.poll(0.2)
              break            
             . . . . . 
              
       def main():
             . . . . .
                  if NOTIFY_TOPIC:
                    for pk, rec in upserts:
                        email = rec.get("email")
                        if email:
                            send_notify(rec["email"], rec.get("id"), rec)
      • main()에서 NOTIFY_TOPIC이 정의되어있으면 send_notify함수를 실행하면서 관련 레코드 인자 외 이메일, ticket_id 값을 전달한다.
      • send_notify 함수는 TOPIC에 전달내용을 payload 디렉터리에 정리해서 producer.produce()호출하여 json 포맷으로 TOPIC(notify.rma) produce한다.
      • producer.poll() 콜백 실행 -> _on_delivery에서 성공/실패 콜백 err 인자로 판별하며 관려 로그를 찍도록 함

5) Notifier Consumer 생성

5-1) Consumer Deployment 생성
apiVersion: apps/v1
kind: Deployment
metadata:
  name: notifier-rma
  namespace: streaming
spec:
  replicas: 1
  selector: { matchLabels: { app: rma-notifier } }
  template:
    metadata: { labels: { app: rma-notifier } }
    spec:
      serviceAccountName: rma-noti-sa
      containers:
      - name: app
        image: ganplank/noti:0.0.8
        imagePullPolicy: IfNotPresent
        env:
        - { name: KAFKA_BOOTSTRAP, value: "doa-kafka-kafka-bootstrap.streaming.svc.cluster.local:9092" }
        - { name: NOTIFY_TOPIC,    value: "notify.rma" }
        - { name: NOTI_GROUP_ID,   value: "notifier-rma" }
        - { name: AUTO_OFFSET_RESET, value: "earliest" }

        # SMTP
        - { name: SMTP_HOST,  value: "smtp.dooray.com" }
        - { name: SMTP_PORT,  value: "465" }
        - { name: SMTP_USER,  valueFrom: { secretKeyRef: { name: smtp-owner, key: username } } }
        - { name: SMTP_PASS,  valueFrom: { secretKeyRef: { name: smtp-owner, key: password } } }
        - { name: MAIL_FROM,  value: "hs.bang@piolink.com" }
        - { name: MAIL_SUBJECT, value: "[RMA] 접수 알림 - {ticket_id}" }
  • 환경변수 내용 및 replicas: 1
  • email 전송을 위한 계정정보는 k8s secret 값을 참조
5-2) 메일 전송 컨슈머
def parse_json(b: Optional[bytes]) -> Dict[str, Any]:
       return json.loads(b.decode("utf-8", errors="replace"))
        
        ....
def safe_commit(consumer: Consumer, msg):
    try:
        consumer.commit(msg, asynchronous=False)
    except Exception as e:
        log(f"[commit-fail] {type(e).__name__}: {e}")
        
        ....

def main():
    conf = {
        "bootstrap.servers": KAFKA_BOOTSTRAP,
        "group.id": NOTI_GROUP_ID,
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
        "session.timeout.ms": 10000,
        "max.poll.interval.ms": 300000,
    }

    consumer = Consumer(conf)
    consumer.subscribe([NOTIFY_TOPIC], on_assign=on_assign, on_revoke=on_revoke)

    try:
        while not stop_flag:
            msg = consumer.poll(POLL_TIMEOUT_SEC)
            if msg is None:
                continue

            tp = f"{msg.topic()}[{msg.partition()}]@{msg.offset()}"
            ts = msg.timestamp()
            raw = msg.value()			   # Partition 에서 가져온 메시지
            
            try:
                payload = parse_json(raw)       # 가져온 Message를 json 포맷을 파이선형태로 디코딩
                handle_record(payload)          # 메일 전송 (실패 시 예외)
                safe_commit(consumer, msg)      # 성공한 경우에만 커밋 (핵심)
  • parse_json() : json 포맷 변환 함수
  • main() : consumer 생성 및 handle_record(payload)을 통해 메시지 전송 후 offset commit
  • handle_record() : 메일 전송 함수
  • safe_commit() : 메일전송 후 동기식으로 offset을 commit
5-3) 메일 메시지 회사 로고 삽입
def build_email(to_email: str, subject: str, html: str, logo_path: str = LOGO_PATH) -> EmailMessage:

	msg = EmailMessage(policy=policy.default)
    msg["Subject"] = subject
    msg["From"] = MAIL_FROM
    msg["To"] = to_email

    msg.set_content("PIOLINK RMA 접수 안내드립니다.") # PLAIN TEXT 파트
    msg.add_alternative(html, subtype="html") # 메일 클라이언트가 subtype = HTML 사용 못하는 경우 PLAIN TEXT 자동 대체

    try:
        with open(logo_path, "rb") as f:
            data = f.read()
        html_part = msg.get_payload()[-1]  
        html_part.add_related(data, maintype="image", subtype="png", cid="logo")
        
    except FileNotFoundError:
        if DEBUG:
            log(f"[warn] logo not found at {logo_path}, skip attaching.")

    return msg
  • EmailMessage(policy=policy.default) : 이메일 메시지작성 클래스로 policy.default는 최신 헤더/인코딩 규칙을 사용하여 작성함
  • add_alternative() : Content-Type: multipart/alternative 사용을 위함이며 이는 여러표현(plain, html) 중 클라이언트쪽에서 렌더링 시 가장 리치한 하나를 선택해서 보여줌
  • 리치한놈 지정은 관련 메소드를 후순위로 코드를 작성하면 된다.
      1. set_content(PLAIN TEXT) 2. add_alternative(HTML) 순서로 코드를 작성하면 HTML을 보여주고 HTML지원되지 않는 클라이언트 경우 1 을 선택한다.
  • 이후 회사로고 이미지 사용을 위해 지정경로에 이미지파일을 읽어와 기존 html에 인라인방식으로 렌더링하고자 multipart/related 으로 msg.get_payload()[-1] text/html 위치에 .add_related 메소드사용하여 cid:logo 부분에 회사로고 이미지가 삽입되도록 함
    • attach()방식은 이미지 파일을 별도 첨부로 전송하거나 외부에서 이미지를 다운로드하는 방식인 반면 인라인방식은 파일첨부 및 이미지 다운로드하지 않고 이미지활용이 가능하기 때문에 보안정책 설정으로 인해 차단되어 이미지가 깨지는 현상을 피할수 있다.
5-4) 메일 메시지 전송
def send_email_ssl_message(msg: EmailMessage):
    attempts, delay = 0, SMTP_BASE_DELAY
    last_exc = None

    while attempts < SMTP_MAX_TRY and not stop_flag:
        try:
            with smtplib.SMTP_SSL(SMTP_HOST, SMTP_PORT, timeout=20) as s:
                # s.set_debuglevel(1)  # 디버그 필요시 주석제거
                if SMTP_USER:
                    s.login(SMTP_USER, SMTP_PASS)
                resp = s.send_message(msg, from_addr=MAIL_FROM, to_addrs=[msg["To"]])
                return
                
        except Exception as e:
            attempts += 1
            last_exc = e
            
            time.sleep(min(delay, SMTP_MAX_DELAY))
            delay = min(delay * 2 + random.random(), SMTP_MAX_DELAY)

    raise last_exc if last_exc else RuntimeError("SMTP send failed")
    
...

def handle_record(rec: Dict[str, Any]):
    to = rec.get("email") or rec.get("record", {}).get("email")
    ticket_id = rec.get("ticket_id") or rec.get("id") or rec.get("record", {}).get("id") or ""
    record = rec.get("record") or rec

    html = html_mail_body(ticket_id, record)
    subject = f"[RMA] 접수 안내 - {ticket_id}"
    msg = build_email(to, subject, html, LOGO_PATH)
    send_email_ssl_message(msg)  # 성공/실패는 상위에서 커밋 제어
  • html_mail_body() : 이메일 전송 내용인 html을 리턴해줌
  • handle_record() : build_email() 함수로 리턴받은 이메일 메시지
  • send_email_ssl_message() : smtp TCP/465포트 메일전송

3. 기타

일전에 github에 .gitignore 설정하지않고 google token값을 push해서 token 사용이 정지되고 재발급받은적이 있다. 이런 경우 git기록에서 제거하고 commit을 되돌리는 작업까지해서 완전히 기록을 제거하는 작업을 하고 다시 push 해야한다.
1) 파일삭제
git rm --cached kafka/upserter/bhs-auth-ea3c2bfa0575.json
echo "kafka/upserter/bhs-auth-ea3c2bfa0575.json" >> .gitignore
2) 이미 포함된 commit을 rewrite
git filter-branch --force --index-filter \ "git rm --cached --ignore-unmatch kafka/upserter/bhs-auth-ea3c2bfa0575.json" \ --prune-empty --tag-name-filter cat -- --all

profile
안녕?

0개의 댓글