Dead Letter Queue

Dead Letter Queue (DLQ) 처리 구성

이 페이지에서는 스트리밍 커넥터(Kafka 고성능, Kinesis 고성능)의 실패 레코드를 버리는 대신 전용 대상(Kafka 토픽/Kinesis 스트림 또는 Snowflake 테이블)으로 라우팅하는 방법을 설명해요. 공통 구성 요소와 DLQ 테이블·권한·원시/구조화 분기·파이프라인을 다룹니다.

출처: Snowflake 문서

본문

Note

이 커넥터는 Snowflake Connector Terms에 의해 규율됩니다.

스트리밍 커넥터(Kafka 고성능, Kinesis 고성능)는 Consume* 프로세서로 메시지를 소비하고 PublishSnowpipeStreaming으로 Snowflake에 전달합니다. 실패는 양쪽에서 발생합니다:

  • 클라이언트 측(Openflow) 실패 — Snowflake에 도달하기 전 파싱·변환에 실패한 레코드. 기본적으로 이것들은 버려집니다(parse-failure 관계가 auto-terminate됨), 그래서 DLQ 처리를 추가하지 않으면 조용히 손실됩니다. 이 토픽은 그러한 실패 레코드를 전용 대상으로 라우팅하는 공통 구성 요소를 설명합니다.

  • 서버 측(Snowpipe Streaming) 실패 — Snowflake에 도달했지만 스키마 불일치 등으로 수집할 수 없는 레코드. 오류 테이블이 활성화되면 서버 측 오류 테이블에 저장됩니다 — Snowpipe Streaming error tables 참고. 오류 테이블은 기본적으로 활성화되지 않으므로, 그것 없이는 이런 실패도 조용히 버려집니다. 여기서 설명하는 DLQ 처리는 클라이언트 측 실패만을 위한 것입니다.

Tip

이 사용자 정의를 손으로 적용할 필요 없습니다. Snowflake CoCo의 Openflow 스킬이 대신 수행할 수 있습니다 — 원하는 변경을 설명하면 이 페이지의 단계에 따라 흐름을 편집합니다. 구성 요소를 수동으로 구성하는 대신 이 스킬을 사용하는 것을 권장합니다.

이것은 공유 참조입니다. 커넥터별 부분 — 소스 프로세서와 DLQ 대상으로 Kafka 토픽/Kinesis 스트림 구성 — 은:

DLQ는 실패 레코드를 두 가지 대상 중 하나로 보낼 수 있습니다:

  • Kafka 토픽 / Kinesis 스트림 — 실패 페이로드를 그대로, 봉투(엔벨로프) 없이 메시징 대상에 재게시. 커넥터별 — Kafka용 Route A / Kinesis용 Route A 참고.

  • Snowflake 테이블 — 실패를 JSON 봉투로 감싸 전용 DLQ 테이블에 삽입(공통 — Route B 참고).

Snowflake 테이블 경로에서는 각 실패 레코드를 균일한 JSON 봉투로 감싸 테이블 열로 조회할 수 있게 합니다. 메시징 대상 경로는 아무것도 감싸지 않습니다 — 원래 실패 페이로드를 그대로 재게시합니다(Kafka/Kinesis 소비자는 원래 바이트를 원함). 따라서 아래 봉투는 Route B에만 적용됩니다. 필드는 제안된 구조입니다 — 필요에 맞게 필드를 추가·제거·이름 변경할 수 있습니다. 봉투를 바꾸면 DLQ 테이블 스키마(Route B 참고)와 raw/structured 분기의 메타데이터 필드를 그것과 일관되게 유지하세요.

Field Description
raw_payload The original, unparseable bytes captured as a string.
structured_payload The parsed JSON record, when available (only for transformation failures).
error_message A short description of why the record failed.
failure_timestamp UTC timestamp of the failure, stored as TIMESTAMP_NTZ.

범위

이 토픽은 이미 설치된 스트리밍 커넥터를 Openflow UI에서 구성 요소를 추가하고 실패 관계를 다시 연결해 그 자리에서 사용자 정의합니다.

DLQ 복잡도는 커넥터가 포함하는 것에 따라 달라집니다:

  • 표준/수정 없는 커넥터는 실패 소스가 Consume* parse failure 관계 하나뿐입니다. 이런 경우 raw-only DLQ를 구축하세요: 파싱할 수 없는 바이트를 raw_payload 문자열로 캡처. 구조화된 처리는 필요 없습니다.

  • 커넥터에 사용자 정의 프로세서나 Custom Transformations 프로세스 그룹이 있다면 그것들은 구조화된(유효한 JSON) 실패를 만들 수 있습니다. 그런 경우에만 선택적 structured_payload 분기를 추가하세요.

범위 밖 — 주요 PublishSnowpipeStreaming(PSS) 전달 실패. PSS는 전달을 재시도하고, Snowflake는 수집할 수 없는 행을 서버 측 오류 테이블에 저장합니다(활성화된 경우) — Snowpipe Streaming error tables 참고. 주요 PSS failure 또는 invalid 관계를 DLQ에 연결하지 마세요. 커넥터가 제공하는 그대로 두세요. failure 관계는 통신 오류를, invalid 관계는 적어도 하나의 레코드가 유효하지 않게 된 FlowFile을 라우팅합니다.

Note

예외. 모든 오류(PSS 전달 실패 포함)를 같은 DLQ 테이블에 넣고 싶다면 주요 PSS failure 관계도 DLQ에 연결할 수 있습니다. 이것은 권장되지 않습니다: 서버 측 오류 테이블이 이미 캡처하는 것을 중복하고 흐름에 부하를 더합니다. 단일 통합 오류 대상이 하드 요구 사항일 때만 사용하세요.

전제 조건

  • Openflow에 배포된 기존 Kafka 고성능 또는 Kinesis 고성능 커넥터가 있어야 합니다.

  • 실패 레코드가 어디로 가야 하는지(Kafka 토픽 / Kinesis 스트림, 또는 Snowflake 테이블) 알고 있어야 합니다.

  • Snowflake 테이블 대상의 경우: execute-as 역할이 대상 테이블에 데이터를 삽입할 수 있어야 합니다(Grants 참고).

  • 메시징 대상(Kafka 토픽 / Kinesis 스트림)의 경우: 커넥터별 페이지를 따르세요.

오류 소스

# Source When Routed to
1 Consume* parse failure Always RAW branch (unparseable bytes)
2 Custom processor / Custom Transformations failure Only if such components exist and you opt into structured handling STRUCTURED branch (or RAW branch if you chose raw-only)
3 Main PublishSnowpipeStreaming failure / invalid Never (out of scope) Left as-is (failure = communication errors; invalid = at least 1 invalid record)

raw-vs-structured 선택은 Snowflake 테이블 경로에만 적용됩니다 — structured_payload는 테이블 열입니다. 스트림 경로는 그에 관계없이 봉투 콘텐츠를 게시합니다.

Note

parse-failure 관계 이름은 커넥터별입니다: parse failure(Kafka, 공백 포함) vs parse.failure(Kinesis, 점 포함). 정확한 이름은 커넥터 페이지를 참고하세요.

공통 설정

아래 공유 구성 요소를 먼저 설정하세요 — DLQ 테이블, 권한, 캡처 분기, 깔때기(funnel), 싱크 실패 처리. 그런 다음 대상 경로를 선택하세요: Route A(Kafka 토픽 / Kinesis 스트림) 또는 Route B(Snowflake 테이블). Snowflake 테이블 경로는 이 모든 구성 요소를 사용합니다. 메시징 경로는 실패 연결만 재사용합니다 — 원래 페이로드를 그대로, 테이블이나 봉투 없이 게시합니다.

테이블 설정

DLQ 테이블이 어디에 있을지 정하세요 — 주요 대상과 같은 데이터베이스·스키마(기본값, 기존 Snowflake Destination Database / Snowflake Destination Schema 파라미터 재사용) 또는 다른 곳.

테이블을 만드세요(원하는 대로 데이터베이스/스키마 조정):

CREATE TABLE pipeline_dlq (
    error_message      VARCHAR,
    failure_timestamp  TIMESTAMP_NTZ,
    raw_payload        VARCHAR,
    structured_payload VARIANT
);

raw-only DLQ에서도 이 전체 스키마를 유지하세요 — structured_payload는 그냥 null로 남습니다. 그렇게 하면 나중에 구조화된 처리를 추가해도 테이블 변경이 필요 없습니다.

권한

execute-as 역할에 다음 권한이 필요합니다:

GRANT USAGE ON DATABASE <db> TO ROLE OPENFLOW_<RUNTIME_NAME>_EXECUTE_AS_RL;
GRANT USAGE ON SCHEMA <db>.<schema> TO ROLE OPENFLOW_<RUNTIME_NAME>_EXECUTE_AS_RL;
GRANT INSERT ON TABLE <db>.<schema>.<table> TO ROLE OPENFLOW_<RUNTIME_NAME>_EXECUTE_AS_RL;

원시 분기 (항상)

비-JSON / 파싱할 수 없는 콘텐츠를 raw_payload 문자열 필드로 캡처합니다. 이것은 표준 커넥터의 유일한 분기이며, 구조화된 분기의 폴백입니다.

Note

한 번에 한 줄이 아니라 전체 페이로드를 캡처하세요. 실패 콘텐츠를 줄 단위로 읽으면(예: No Match Behavior = raw-line인 GrokReader) 줄마다 레코드 하나를 배출하므로, 여러 줄 페이로드(예쁘게 출력된 JSON, 여러 줄 텍스트)가 많은 가짜 DLQ 행으로 쪼개집니다. 대신 ExtractText → UpdateAttribute → AttributesToJSON을 사용해 전체 FlowFile 콘텐츠를 단일 raw_payload 값으로 캡처하세요. 줄바꿈과 관계없이 모든 비-JSON 콘텐츠에 동작합니다.

1단계 — ExtractText 프로세서 추가("Capture Whole Payload"). DOTALL을 켜서 새 줄이 포함되도록 전체 콘텐츠를 raw_payload 속성으로 복사합니다:

Property Value
Enable DOTALL Mode true
Maximum Buffer Size 1 MB
Maximum Capture Group Length 1048576
raw_payload (dynamic property) (?s)(.*)

Note

바이너리 또는 매우 큰 페이로드는 전체 콘텐츠를 읽고 {"raw_payload": <json-escaped-content>, ...}을 직접 쓰는 작은 ExecuteGroovyScript를 대신 사용하세요.

2단계 — UpdateAttribute 프로세서 추가("Add Raw DLQ Metadata"):

Property Value
error_message raw payload (non-JSON) — on Kinesis you can use ${record.error.message} instead (see the Kinesis page)
failure_timestamp ${now():format('yyyy-MM-dd HH:mm:ss.SSS', 'UTC')}

3단계 — AttributesToJSON 프로세서 추가("Build Raw DLQ Envelope"). 봉투를 FlowFile 콘텐츠에 씁니다(값은 자동으로 JSON 이스케이프됨):

Property Value
Attributes List raw_payload,error_message,failure_timestamp
Destination flowfile-content
Include Core Attributes false

4단계 — PublishSnowpipeStreaming 프로세서 추가("Ingest Failed Records into DLQ Table"). 기존 PSS auth / web-client 서비스를 재사용합니다. DLQ 테이블을 가리키게 하세요:

Property Value
Destination Type TABLE
Database The DLQ table's database.
Schema The DLQ table's schema.
Table The DLQ table name (for example, pipeline_dlq).
Channel Group ${hostname(false)}.dlq
Authentication Strategy SNOWFLAKE_MANAGED
Connection Strategy STANDARD
Web Client Service Provider The existing web-client service.

success, empty를 auto-terminate하세요. failure, invalid를 DLQ 싱크 실패 처리로 라우팅합니다.

구조화된 분기 (조건부)

사용자 정의 구성 요소가 있고 파싱된 JSON을 structured_payload VARIANT 열에 보존하려는 경우에만 이 분기를 만드세요.

1단계 — JoltTransformRecord 프로세서 추가("Move Payload to structured_payload + Add Metadata"). 단일 프로세서가 레코드를 structured_payload 아래로 옮기고 봉투 메타데이터를 추가합니다 — Jolt Specification 속성은 Expression Language를 지원합니다:

Property Value
Record Reader The existing JsonTreeReader.
Record Writer The existing JsonRecordSetWriter.
Jolt Transform jolt-transform-chain
Jolt Specification The chain spec below.
[
  {"operation": "shift",   "spec": {"*": "structured_payload.&"}},
  {"operation": "default", "spec": {
    "error_message": "pipeline failure",
    "failure_timestamp": "${now():format('yyyy-MM-dd HH:mm:ss.SSS','UTC')}"
  }}
]

(Kinesis에서는 error_message을 ${record.error.message}로 설정 — Kinesis page 참고.) 그러면 raw 분기와 같은 봉투 필드(error_message, failure_timestamp, 그리고 structured_payload)가 만들어집니다. original을 auto-terminate하세요. failure를 raw 깔때기로 라우팅합니다(폴백: 구조화 실패 시에도 바이트를 캡처).

JoltTransformRecord success를 raw 분기에서 구성한 같은 PublishSnowpipeStreaming 프로세서("Ingest Failed Records into DLQ Table")로 라우팅하세요. 별도 PSS는 필요 없습니다 — 두 분기가 같은 봉투 스키마를 만들고 같은 채널 그룹(${hostname(false)}.dlq)을 공유합니다.

깔때기와 연결

여러 실패 소스가 깔끔하게 수렴하도록 깔때기(funnel) 를 병합 지점으로 사용하세요:

  • RAW 깔때기 — 모든 파싱 불가 / ser-de 실패 소스(Consume* parse-failure 관계, 그리고 구조화 분기 실패)가 여기서 수렴해 ExtractText → UpdateAttribute → AttributesToJSON → 단일 PublishSnowpipeStreaming("Ingest Failed Records into DLQ Table")으로 흐릅니다.

  • JSON 깔때기(구조화 분기가 있을 때만) — 사용자 정의 / 변환 실패가 여기서 수렴해 JoltTransformRecord → 같은 단일 PublishSnowpipeStreaming으로 흐릅니다.

raw-only 커넥터라면 JSON 깔때기와 구조화 분기를 완전히 생략하고, 모든 소스가 RAW 깔때기로 가게 하세요.

DLQ 싱크 실패 처리

DLQ 게시/삽입조차 실패하면 레코드를 잃지 마세요:

  1. DLQ 싱크의 failure, invalid 관계(PublishSnowpipeStreaming, PublishKafka, 또는 PublishKinesis)를 LogAttribute 프로세서("Log DLQ Ingestion Error", Log Level = error, "Failed to ingest data into the DLQ" 같은 프리픽스)로 라우팅합니다.

  2. LogAttribute success를 검사를 위해 전달되지 못한 레코드가 쌓이는 parking-lot 깔때기로 라우팅합니다.

Route A — Kafka 토픽 / Kinesis 스트림

이 경로는 실패 페이로드를 그대로, 봉투 없이 메시징 대상에 재게시합니다 — 공통 설정의 DLQ 테이블이나 raw/structured 분기를 사용하지 않습니다. 게시자와 그것이 재사용하는 연결/자격 증명은 커넥터별이므로 단계는 커넥터 페이지에 있습니다:

Route B — Snowflake 테이블

이 경로는 공통 설정의 구성 요소(DLQ 테이블, 권한, raw/structured 분기, 깔때기)를 PublishSnowpipeStreaming을 종단 싱크로 조립합니다. 두 커넥터에 동일합니다.

데이터 타입과 JSONL 처리

두 싱크 모두 같은 봉투를 소비합니다: PublishSnowpipeStreaming은 JSONL(줄마다 JSON 객체 하나)을 기대하고, 스트림 게시자는 그 같은 JSONL 콘텐츠를 메시지 값으로 보냅니다. 실패는 서로 다른 형태로 오므로, 유효한 레코드를 만드는 리더로 라우팅하세요:

  • 파싱 불가 / 비-JSON(parse failure): raw 분기의 ExtractText(DOTALL)가 전체 페이로드 — 여러 줄을 포함한 어떤 텍스트든 — 단일 raw_payload 필드로 캡처합니다. (줄 기반 리더(No Match Behavior = raw-line인 GrokReader 등)는 레코드당 단일 텍스트 줄만 감싸므로 여러 줄 페이로드를 쪼갭니다; 전체 콘텐츠 캡처를 쓰세요.) 이것이 raw 분기를 비정상 입력에 대해 강건하게 만드는 이유입니다.

  • 유효한 JSON(다운스트림/변환 실패): 구조화 분기의 Jolt shift {"*":"structured_payload.&"}가 파싱된 객체를 structured_payload 아래에 중첩합니다. 어떤 실패에서든 raw 깔때기로 폴백하므로 잃는 것이 없습니다.

Note

흐름을 따라 데이터 타입 변경을 주의하세요. 메시지는 소스에서 한 타입(예: CSV/TSV 또는 꺠진 줄)이고 다운스트림 매핑 후에는 다른 타입일 수 있습니다. 실패 지점의 콘텐츠에 맞게 캡처하세요 — 원시 CSV/TSV를 JSON 리더에 넣지 마세요. 전체 콘텐츠 raw 분기는 모든 비-JSON 콘텐츠를 안전하게 처리합니다.

검증

연결 후 먼저 모든 컨트롤러 서비스를 활성화하고(비활성 서비스 참조 프로세서는 INVALID로 표시), 그다음 프로세스 그룹을 검증하세요. 흐름을 시작하기 전에 검증 실패를 고치세요.

문제 해결

Symptom Likely cause
CREATE TABLE / insert denied Missing grants on the DLQ schema (see Grants).
Records pile up in the parking-lot funnel DLQ table/schema mismatch or wrong database/schema/table parameters — inspect the LogAttribute error output.

커넥터별 증상(parse-failure 관계 이름, 토픽/스트림 게시)은 Kafka 및 Kinesis 문제 해결 표를 참고하세요.

더 알아보기 (Learn more)