Kafka용 Openflow 커넥터: DLQ 처리 구성

Kafka용 Openflow 커넥터: DLQ 처리 구성

이 주제는 Kafka 고성능 커넥터에서 Dead Letter Queue(DLQ) 메시지의 대상으로 Kafka 토픽을 구성하는 방법과, DLQ 처리의 다른 Kafka 특화 부분(소스 프로세서, 파싱 실패 관계, 연결 재사용)을 설명해요.

출처: Snowflake 문서

본문

참고: 이 커넥터는 Snowflake Connector Terms가 적용돼요.

이 주제는 Kafka 최종 대상으로 Kafka high-performance 커넥터에서 Dead Letter Queue(DLQ) 메시지를 구성하는 방법과, DLQ 처리의 다른 Kafka 특화 부분(소스 프로세서, 파싱 실패(parse-failure) 관계, 연결 재사용)을 설명해요.

DLQ 처리는 스트리밍 커넥터들 사이에서 공유돼요. 공통 개념(실패 봉투(envvelope), Snowflake 테이블 경로, raw/structured 분기, funnel, DLQ 싱크 실패 처리)을 먼저 일반 가이드인 Configuring Dead Letter Queue (DLQ) handling에서 읽어 보세요. 이 페이지는 Kafka에 특화된 부분만 다뤄요.

팁: 이 커스터마이즈를 직접 손으로 적용할 필요는 없어요. Snowflake CoCo의 Openflow 스킬이 대신 수행해 줄 수 있어요 — 원하는 변경을 설명하면 이 페이지의 단계에 따라 플로우를 편집해요. 구성 요소를 수동으로 구성하는 것보다 스킬을 사용하는 것을 권장해요.

커넥터 근거(Connector grounding)

항목 Kafka 고성능
소스 프로세서 ConsumeKafka
파싱 실패 관계 parse failure
재사용할 연결/자격 증명 Kafka3ConnectionService
스트림 경로 게시자 PublishKafka
레코드 리더/라이터 JsonTreeReader / JsonRecordSetWriter
대상 프로세서 PublishSnowpipeStreaming

파싱 실패를 DLQ로 라우팅하기

새 커넥터에서는 ConsumeKafka의 parse failure 관계가 자동 종료(auto-terminated)돼요. 자동 종료를 제거하고 공통 가이드에 설명된 RAW funnel에 parse failure를 연결해요.

DLQ 메시지 대상으로서의 Kafka 토픽

이 경로를 사용해 실패한 레코드를 Kafka 토픽으로 다시 게시해요. 원래 실패한 페이로드를 있는 그대로 게시하세요 — 봉투(envelope)도 레코드 래핑도 없어요. 실패 소스를 PublishKafka 싱크에 직접 연결해요. 봉투(raw_payload / structured_payload)는 Snowflake 테이블 경로 전용인데, Kafka 컨슈머는 원래 바이트를 원하기 때문이에요.

오류 컨텍스트는 메시지 본문 안이 아니라 Kafka 헤더로 대역 외(out-of-band) 전달돼요(FlowFile Attribute Header Pattern을 통해).

참고: 연결 재사용: DLQ 게시자는 ConsumeKafka와 같은 Kafka3ConnectionService를 재사용해요 — 즉 같은 클러스터예요. DLQ 토픽이 다른 Kafka 클러스터에 있다면, 해당 클러스터에 맞게 구성된 별도의 Kafka3ConnectionService를 만들고 PublishKafka가 그것을 가리키게 해요. 같은 커넥터가 DLQ 메시지를 어떤 Kafka 클러스터와 어떤 Snowflake 환경에도 쓸 수 있어요. 그렇지 않으면 기존 연결이 재사용돼요.

1단계: PublishKafka 프로세서 만들기

  1. 커넥터의 프로세스 그룹에 PublishKafka 프로세서를 추가해요.
  2. 다음 속성을 설정해요:
속성 값
Topic Name DLQ 토픽 이름
Kafka Connection Service ConsumeKafka가 사용하는 것과 같은 Kafka3ConnectionService
Failure Strategy Route to Failure
FlowFile Attribute Header Pattern kafka\\..*

2단계: 실패 소스를 게시자에 연결하기

  • 실패 소스(parse failure 관계와 변환/오류 관계들)를 이 게시자에 직접 연결해요 — 스트림 경로에는 raw/structured 분기가 만들어지지 않아요.
  • 게시자의 failure, invalid 관계를 DLQ sink failure handling로 라우팅해요. 제한된(bounded) failure 재시도(예: 재시도 횟수 3)를 사용해서 일시적인 브로커 문제는 복구되지만 영구적인 실패는 여전히 파킹 롯(parking-lot)에 도달하게 해요. 사실상 무한 재시도(예: 9999)는 사용하지 마세요.

DLQ 메시지 대상으로서의 Snowflake 테이블

두 커넥터 모두에서 동일해요. 공통 가이드의 Route B — Snowflake table을 참고해요.

문제 해결(Troubleshooting)

증상 가능한 원인
DLQ 게시자가 잘못된 클러스터에 씀 PublishKafka가 ConsumeKafka의 Kafka3ConnectionService를 재사용함 — 다른 클러스터는 별도의 연결 서비스가 필요함

공유 증상(raw 분기, structured_payload, 권한(grants), 파킹 롯 funnel)은 공통 문제 해결 표를 참고해요.

더 알아보기 (Learn more)