Kinesis 데이터 스트림용 Openflow 커넥터: DLQ 처리 구성
Kinesis 데이터 스트림용 Openflow 커넥터: DLQ 처리 구성
이 문서에서는 Kinesis 고성능 커넥터에서 Dead Letter Queue(DLQ) 메시지의 목적지로 Kinesis 스트림을 구성하는 방법과, DLQ 처리에서 Kinesis에 특화된 나머지 부분(소스 프로세서, 파싱 실패 관계, 자격 증명 재사용)을 설명합니다.
출처: Snowflake 문서
본문
DLQ 처리는 스트리밍 커넥터 간에 공유됩니다. 공통 개념은 먼저 일반 가이드인 [Dead Letter Queue(DLQ) 처리 구성]을 읽어 주세요. 여기에는 실패 봉투(failure envelope), Snowflake 테이블 경로, 원시/구조화(raw/structured) 분기, 퍼널(funnel), DLQ 싱크 실패 처리 등이 포함됩니다. 이 페이지에서는 Kinesis에만 해당하는 내용만 다룹니다.
팁: 이 커스터마이징을 직접 적용할 필요는 없습니다. Snowflake CoCo의 Openflow 스킬이 이 작업을 대신 수행해 줍니다. 원하는 변경을 설명하면 이 페이지의 단계를 따라 플로우를 수정해 줍니다. 구성 요소를 수동으로 구성하는 대신 스킬을 사용하는 것을 권장합니다.
커넥터 접지(Connector grounding)
| 항목 | Kinesis high-performance |
|---|---|
| 소스 프로세서 | ConsumeKinesis |
| 파싱 실패 관계 | parse.failure |
| 재사용할 연결/자격 증명 | AWSCredentialsProviderControllerService + Region + Stream Name |
| 스트림 경로 퍼블리셔 | PutKinesisStream |
| 레코드 리더/라이터 | JsonTreeReader / JsonRecordSetWriter |
| 목적지 프로세서 | PublishSnowpipeStreaming |
파싱 실패를 DLQ로 라우팅
새 커넥터에서 ConsumeKinesis의 parse.failure 관계는 자동 종료(auto-terminate)됩니다. 자동 종료를 제거하고 parse.failure를 일반 가이드에 설명된 RAW 퍼널에 연결하세요.
팁: 오류 이유를 캡처하세요.
ConsumeKinesis는 파싱/serde 실패 시record.error.messageFlowFile 속성을 기록합니다. 원시/구조화 분기 메타데이터의error_message필드에는${record.error.message}를 사용하세요(카프카 커넥터에는 이에 상응하는 속성이 없습니다).
DLQ 메시지의 목적지로서의 Kinesis 스트림
이 경로를 사용하여 실패한 레코드를 다시 Kinesis 스트림으로 게시하세요. 원래 실패한 페이로드를 그대로 게시합니다. 봉투도 레코드 래핑도 없습니다. 실패 소스를 PutKinesisStream 싱크에 직접 연결하세요. 봉투(raw_payload/structured_payload)는 Snowflake 테이블 경로 전용입니다. 스트림 소비자는 원래 바이트를 원하기 때문입니다.
참고: 자격 증명 재사용 - DLQ 퍼블리셔는
ConsumeKinesis와 동일한AWSCredentialsProviderControllerService+ Region, 즉 동일한 AWS 계정/리전을 재사용합니다. DLQ 스트림이 다른 계정이나 리전에 있다면 적절한 자격 증명/리전으로 퍼블리셔를 구성하세요(완전히 분리된 환경이면 별도의 커넥터를 사용하세요).
1단계: PutKinesisStream 프로세서 만들기
- 커넥터의 프로세스 그룹에 PutKinesisStream 프로세서를 추가하세요.
- 다음 속성을 설정하세요.
| 속성 | 값 |
|---|---|
| Stream Name | DLQ 스트림 이름 |
| AWS Credentials Provider Service | ConsumeKinesis가 사용하는 것과 동일한 AWSCredentialsProviderControllerService |
| Region | ConsumeKinesis와 동일한 리전 |
참고:
PutKinesisStream은 FlowFile 콘텐츠 전체를 단일 Kinesis 메시지로 게시합니다. FlowFile에 여러 레코드가 들어 있다면(예: 한 줄에 레코드 하나씩인 NDJSON)PutKinesisStream전에SplitText프로세서를 사용해 줄 단위로 FlowFile을 나누어 주세요.
2단계: 실패 소스를 퍼블리셔에 연결
- 실패 소스(parse.failure 관계 및 변환/오류 관계)를 이 퍼블리셔에 직접 연결하세요. 스트림 경로에는 원시/구조화 분기를 만들지 않습니다.
- 퍼블리셔의 failure, invalid 관계를 DLQ 싱크 실패 처리로 라우팅하세요. 제한된 실패 재시도(예: 재시도 횟수 3)를 사용해 일시적인 스트림 문제는 복구되되 영구적인 실패는 여전히 parking-lot에 도달하게 하세요. 사실상 무한 재시도(예: 9999)는 사용하지 마세요.
DLQ 메시지의 목적지로서의 Snowflake 테이블
두 커넥터 모두 동일합니다. 일반 가이드의 Route B(Snowflake 테이블)를 참조하세요.
문제 해결
| 증상 | 가능한 원인 |
|---|---|
| DLQ 퍼블리셔가 잘못된 스트림/계정에 기록 | PutKinesisStream이 ConsumeKinesis의 AWSCredentialsProviderControllerService + Region을 재사용함. 다른 계정/리전은 적절한 자격 증명과 리전이 필요함. |
공통 증상(원시 분기, structured_payload, grants, parking-lot 퍼널)은 공통 문제 해결 표를 참조하세요.