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.message FlowFile 속성을 기록합니다. 원시/구조화 분기 메타데이터의 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 퍼널)은 공통 문제 해결 표를 참조하세요.

더 알아보기 (Learn more)