
notify.rma KafkaTopic으로 메시지를 전송하면 Notifier Consumer가 Topic(Partition)에 메시지를 읽고 메일과 슬랙 메시지 전송 작업을한다.


nohup sudo kubectl -n rma port-forward --address 0.0.0.0 svc/rma-web-svc 80:80 --kubeconfig $HOME/.kube/config &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
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)
producer.poll() 콜백 실행 -> _on_delivery에서 성공/실패 콜백 err 인자로 판별하며 관려 로그를 찍도록 함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}" }
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 commithandle_record() : 메일 전송 함수safe_commit() : 메일전송 후 동기식으로 offset을 commitdef 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) 중 클라이언트쪽에서 렌더링 시 가장 리치한 하나를 선택해서 보여줌msg.get_payload()[-1] text/html 위치에 .add_related 메소드사용하여 cid:logo 부분에 회사로고 이미지가 삽입되도록 함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포트 메일전송일전에 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