사용자 정의 변환 구성

사용자 정의 변환 구성 (Custom transformations)

이 페이지에서는 스트리밍 커넥터(Kafka 고성능, Kinesis 고성능)의 소스와 대상 사이에 사용자 정의 처리를 끼워 넣는 방법을 설명해요. 필터링·필드 매핑(평탄화/이름 변경/제거)·토픽-테이블 매핑·콘텐츠 기반 라우팅·기본값·사용자 정의 Groovy 스크립트를 Custom Transformations 프로세스 그룹으로 추가합니다.

출처: Snowflake 문서

본문

Note

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

스트리밍 커넥터(Kafka 고성능, Kinesis 고성능)는 Consume* 프로세서로 메시지를 소비하고 PublishSnowpipeStreaming으로 Snowflake에 전달합니다. 기본적으로 소스는 PublishSnowpipeStreaming에 직접 연결됩니다. 이 토픽은 Openflow UI에서 Custom Transformations 프로세스 그룹을 추가해 소스와 대상 사이에 사용자 정의 처리(필터링, 필드 매핑(평탄화/이름 변경/제거), 토픽-테이블 매핑, 여러 테이블로의 콘텐츠 기반 라우팅, 기본값, 사용자 정의 Groovy 스크립트)를 끼워 넣는 방법을 설명합니다. 이미 배포된 스트리밍 커넥터에 적용합니다. 아직 없다면 먼저 Kafka 또는 Kinesis 커넥터를 설정하세요.

스트리밍 커넥터는 at-least-once(ALO) 전달을 사용합니다. 사용자 정의 처리는 그 경로 안에 직접 놓이므로 제한 규칙을 지켜야 합니다. 지키지 않으면 중복, 순서 어긋남, 데이터 손실(레코드를 명시적으로 버린 경우)이 발생할 수 있습니다. 프로세서를 추가하기 전에 제한 규칙을 읽으세요.

Tip

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

시작하려면 Snowflake CoCo CLI를 설치·연결한 뒤 번들된 openflow 스킬에 변경을 요청하세요.

메시지 데이터 타입 전환(JSON → Avro / Protobuf)은:

실패 레코드를 버리는 대신 전용 대상으로 라우팅하는 것은:

범위

이 토픽은 이미 설치된 스트리밍 커넥터를 Openflow UI에서 프로세스 그룹과 프로세서를 추가하고 소스-대상 경로를 다시 연결해 그 자리에서 사용자 정의합니다.

이 토픽이 다루는 것:

  • 가독성을 위해 선택적으로 Custom Transformations 프로세스 그룹 삽입(프로세서를 커넥터 캔버스에 직접 추가할 수도 있음).

  • 안전한 in-flight 처리를 위한 제한 규칙.

  • 커넥터의 기존 JsonTreeReader / JsonRecordSetWriter 재사용과 스키마 캐시 추가.

  • 직렬화/역직렬화(ser/de) 패스를 최소화하는 파이프라인 계획.

  • 일반적인 변환 패턴과 이를 구현하는 프로세서들.

범위 밖 — 브로커 인증, 데이터 타입 전환, Snowflake Private Key 인증. See also 링크를 참고하세요.

전제 조건

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

  • 커넥터가 소스 프로세서(예: ConsumeKafka)를 PublishSnowpipeStreaming에 직접 연결하고 있어야 합니다.

  • 어떤 변환이 필요하고 레코드를 어떤 테이블(들)에 쓸지 알고 있어야 합니다.

  • execute-as 역할이 라우팅에 사용되는 대상 테이블을 만들 수 있어야 합니다(커넥터의 대상 권한 참고).

아키텍처

소스 프로세서와 PublishSnowpipeStreaming 사이에 변환 프로세서를 끼워 넣으세요 — 즉 소스에서 PublishSnowpipeStreaming으로의 직접 연결을 제거하고 그 사이에 프로세서를 연결합니다.

Note

Custom Transformations 프로세스 그룹은 선택 사항입니다. 변환 프로세서를 커넥터 캔버스에 소스 프로세서와 PublishSnowpipeStreaming 사이에 직접 추가할 수 있습니다. Custom Transformations라는 전용 프로세스 그룹(Input Port와 Output Port 포함) 안에 묶는 것은 순전히 가독성·유지보수성을 위한 것입니다 — 사용자 정의를 커넥터의 내장 구성 요소와 시각적으로 분리해 찾고·이름 붙이고·이해하기 쉽게 만듭니다. 기능적으로 달라지는 것은 없습니다: 프로세서, 제한 규칙, FIFO 연결, 리더/라이터는 어느 쪽이든 동일합니다. 나머지 토픽은 그룹화된 배치를 설명합니다. 그룹을 건너뛴다면 Input/Output Port 단계는 무시하고 소스 프로세서를 첫 번째 변환 프로세서에, 마지막 프로세서를 PublishSnowpipeStreaming에 직접 연결하세요.

권장(그룹화된) 배치. Input Port가 소스에서 메시지를 받고, 그룹 안에서 변환 프로세서가 실행되며, Output Port가 결과를 PublishSnowpipeStreaming에 공급합니다:

ConsumeKafka (or ConsumeKinesis)
    | (success)
    v
[Custom Transformations Input]   <- Input Port
    |
    ... transformation processors ...
    |
[Custom Transformations Output]  <- Output Port
    |
    v
PublishSnowpipeStreaming

1단계: 프로세스 그룹 생성(선택)

이 단계(그리고 2단계)는 그룹화된 배치를 원할 때만 필요합니다. 변환 프로세서를 커넥터 캔버스에 직접 추가하려면 아래 1.1–1.2를 수행한 뒤 3단계로 건너뛰고 소스 프로세서를 첫 번째 변환 프로세서에 직접 연결하세요.

  1. Openflow UI에서 커넥터의 프로세스 그룹을 엽니다.

  2. 소스 프로세서와 PublishSnowpipeStreaming 사이의 기존 연결을 제거합니다.

  3. 새 Process Group을 캔버스로 끌어와 Custom Transformations로 이름을 지정합니다.

2단계: 입력·출력 포트 추가(선택)

1단계에서 프로세스 그룹을 만든 경우에만. Custom Transformations 그룹 안에서:

  1. Custom Transformations Input이라는 Input Port를 추가합니다.

  2. Custom Transformations Output이라는 Output Port를 추가합니다.

3단계: 흐름에 처리 연결

  1. 소스 프로세서의 success 관계를 Custom Transformations Input에 연결합니다(그룹을 건너뛰었다면 첫 번째 변환 프로세서에 직접).

  2. Custom Transformations Output을 PublishSnowpipeStreaming에 연결합니다(그룹을 건너뛰었다면 마지막 변환 프로세서를 PublishSnowpipeStreaming에).

  3. 모든 연결(이 두 외부 연결과 이후 추가하는 모든 내부 연결)을 FirstInFirstOut 우선순위자로 구성합니다. 연결 구성 참고.

제한 규칙

Warning

이 규칙은 반드시 지켜야 합니다. 위반하면 데이터 손실, 중복, 성능 저하가 발생할 수 있습니다.

  1. ser/de 연산 최소화. 프로세서를 만들기 전에 전체 파이프라인을 계획하세요. 필터링·이름 변경·기본값을 가능한 적은 수의 레코드 인식 프로세서에 결합하세요. 유용한 값(예: 라우팅 필드)은 일찍 속성으로 추출해 다운스트림 단계가 콘텐츠 인식 프로세서 대신 속성 전용 프로세서(RouteOnAttribute, UpdateAttribute)를 쓰게 하세요.

  2. 가능하면 기존 리더/라이터 재사용. 커넥터는 이미 부모 프로세스 그룹 수준에서 JsonTreeReader(스키마 추론)와 JsonRecordSetWriter를 정의하며, 이들은 프로세서에 보입니다. 프로세서가 같은 데이터 타입의 콘텐츠를 읽거나 쓰면 새로 만들기보다 기존 서비스를 재사용하세요 — 흐름을 단순하고 일관되게 유지합니다. Reader and writer setup 참고.

  3. FlowFile당 하나의 테이블. 단일 FlowFile은 ONE 테이블의 레코드만 포함할 수 있습니다. 라우팅 값이 이미 속성으로 있다면 PartitionRecord보다 RouteOnAttribute를 선호하세요(ser/de 패스를 아낍니다).

  4. 연결에 FIFO 우선순위자 사용. ALL 연결(내부·외부)에 FirstInFirstOut 우선순위자를 적용하세요. 각 파티션 안에서 레코드 도착 순서를 유지하고 예상치 못한 인터리빙을 피합니다. 처리량에는 영향이 없습니다.

  5. 같은 테이블로의 배열 폭발 없음. 같은 테이블로 향하는 여러 레코드로 배열을 폭발시키는 것은 허용되지 않습니다(데이터 손실이나 중복을 유발할 수 있음). 배열 폭발은 결과 레코드가 다른 테이블로 라우팅될 때만 유효합니다.

  6. 구조는 스키마 추론에 맡기기. 기존 JsonTreeReader(+ VolatileSchemaCache)와 JsonRecordSetWriter를 사용하면 구조 변경은 자동으로 처리됩니다 — 수동으로 스키마를 정의할 필요가 없습니다.

리더/라이터 설정

FlowFile 콘텐츠에 접근하는 변환 프로세서는 Record Reader와 Record Writer가 필요합니다. 커넥터는 이미 부모 프로세스 그룹 수준에서 JsonTreeReader와 JsonRecordSetWriter를 정의하며, 이들은 Custom Transformations 그룹 안에 보입니다. 프로세서가 같은 데이터 타입으로 작업하면 새로 만들기보다 기존 서비스를 재사용하세요.

리더에 스키마 캐시 추가(권장)

기존 JsonTreeReader는 스키마 추론을 사용합니다. 매 FlowFile마다 스키마를 다시 추론하지 않도록 VolatileSchemaCache를 추가하고 리더가 이를 가리키게 하세요.

  1. 커넥터 프로세스 그룹 수준에서 Configure > Controller Services로 이동합니다.

  2. VolatileSchemaCache 서비스를 추가합니다.

Property Value
Maximum Cache Size 100 (default is usually sufficient)
  1. VolatileSchemaCache 서비스를 Enable합니다.

  2. JsonTreeReader를 Disable합니다(편집 전에 필요).

  3. JsonTreeReader를 편집해 설정합니다:

Property Value
Schema Inference Cache The VolatileSchemaCache created above.
  1. JsonTreeReader를 다시 Enable합니다.

이렇게 하면 추론된 스키마를 캐시하고 같은 구조의 메시지에 재사용해 성능을 개선합니다.

변환 계획

프로세서를 만들기 전에 전체 파이프라인을 계획하세요(규칙 1). 목표는 가능한 가장 적은 ser/de 패스입니다. 계획 없이 프로세서를 하나씩 만들면 거의 항상 낭비적인 파이프라인이 됩니다.

계획 체크리스트

  1. 필요한 모든 변환을 나열(필터, 이름 변경, 기본값, 라우팅 등).

  2. 어느 것이 속성만으로 실행될 수 있는지 식별(ser/de 0회). 속성 전용 프로세서는 FlowFile이 이미 해당 속성을 지닐 때만 동작합니다(예: ConsumeKafka가 설정하는 kafka.topic, ConsumeKinesis가 설정하는 aws.kinesis.stream.name, 또는 업스트림 PartitionRecord가 설정한 속성):

  • 소스(토픽/스트림) 또는 기존 속성으로 라우팅 → UpdateAttribute로 대상 테이블 이름을 속성(예: table.name)에 도출. PublishSnowpipeStreaming이 이를 읽어 대상 테이블을 선택합니다. 자세한 내용은 토픽-테이블 매핑과 여러 테이블로의 콘텐츠 기반 라우팅을 참고하세요.

  • 속성 기반 필터링 → RouteOnAttribute.

  1. 콘텐츠 변환을 가장 적은 수의 콘텐츠 인식 프로세서로 결합:
  • 필터 + 이름 변경 + 기본값 → ONE QueryRecord(SQL SELECT 별칭, COALESCE, WHERE) 또는 ONE JoltTransformRecord(Chain 사양).
  1. 다중 테이블 라우팅에 PartitionRecord 사용. FlowFile에는 라우팅 필드 값이 다른 레코드가 많이 있을 수 있습니다. 파티셔닝은 FlowFile을 분할해 각 출력 FlowFile이 한 값의 레코드만 갖게 하고, 그 값을 FlowFile 속성으로 설정합니다.

  2. 작업 순서: 콘텐츠 변환 먼저(전체 FlowFile에 대해 한 패스) → 파티션(값별 FlowFile로 분할) → UpdateAttribute로 파티션된 값에서 대상 테이블 이름 속성 설정(ser/de 0회 — PublishSnowpipeStreaming이 이 속성으로 대상 테이블 선택) → 필터링이 필요할 때만 RouteOnAttribute(ser/de 0회). 테이블 이름 속성은 여러 테이블로의 콘텐츠 기반 라우팅에서 설명합니다.

가독성 vs 최적화 파이프라인 선택

표준 프로세서 파이프라인을 계획한 뒤 ser/de 패스 수를 셉니다:

  • 1 패스 — 최적화할 것이 없으므로 그대로 구현.

  • 2 패스 이상 — 단일 ExecuteGroovyScript가 모든 콘텐츠 연산을 1 패스로 통합할 수 있습니다. 유지보수성(표준 프로세서는 읽고 수정하기 쉬움)과 성능(성능 중시 워크로드에서 Groovy 스크립트 하나가 더 빠를 수 있음 — 최적화 전에 테스트로 가정을 확인하세요, 다만 변경이 어려움)을 저울질하세요. 최대 성능이 특별히 필요한 게 아니면 가독성 있는 옵션을 선택하세요.

예시 — "타임스탬프로 필터링, 필드 이름 변경, 기본값 추가, 필드 값으로 테이블 라우팅":

가독성(ser/de 2 패스):

Input Port
  -> QueryRecord    (ser/de #1: rename via SELECT aliases, defaults via COALESCE, filter via WHERE)
  -> PartitionRecord (ser/de #2: split by routing field -> sets attribute)
  -> UpdateAttribute (zero ser/de: table.name = ${routing-field})
  -> RouteOnAttribute (zero ser/de: drop unwanted values, auto-terminate unmatched)
  -> Output Port

최적화(ser/de 1 패스):

Input Port
  -> ExecuteGroovyScript (ser/de #1: rename + defaults + filter + partition by routing field)
  -> UpdateAttribute (zero ser/de: table.name = ${routing-field})
  -> RouteOnAttribute (zero ser/de: drop unwanted values)
  -> Output Port

PublishSnowpipeStreaming은 다중 테이블 라우팅을 기본적으로 처리합니다 — Table 속성을 ${table.name}(또는 라우팅 속성 직접)으로 설정하세요. RouteOnAttribute는 필터링에만 필요하며, 라우팅만을 위해 필요한 것은 아닙니다.

안티 패턴(피해야 할 것): 필터/이름 변경/파티션 각각에 별도 프로세서(2면 충분한데 3 ser/de 패스), 또는 라우트 값마다 별도 UpdateAttribute. 가능하면 결합하고, 파티션된 속성으로 구동되는 단일 UpdateAttribute를 쓰세요.

연결 구성

모든 연결 — Custom Transformations 그룹 안팎, 그리고 그룹 안 프로세서 사이 — 반드시 다음과 같이 구성해야 합니다:

  1. FirstInFirstOut 우선순위자 — 각 파티션 안에서 레코드 도착 순서를 유지(규칙 4).

  2. Back pressure — 특별한 이유가 없으면 기본값 유지.

우선순위자를 설정하려면 연결을 편집하고 Settings 탭을 열고 Prioritizers 아래에 FirstInFirstOutPrioritizer를 추가하세요.

이것은 다음에 적용됩니다:

  • 소스 프로세서 → Custom Transformations Input

  • Custom Transformations Output → PublishSnowpipeStreaming

  • 그룹 안 프로세서 사이의 모든 연결

Note

PublishSnowpipeStreaming 앞의 연결 큐: 연결 큐 크기 한도를 5 GB로 설정하세요(기본 1 GB 대신). PublishSnowpipeStreaming은 Snowflake가 확인할 때까지 FlowFile을 큐에 보관합니다. 기본 1 GB 한도에서는 정상 부하에서 백프레셔가 너무 일찍 발동합니다.

변환 패턴

필요에 맞는 패턴을 고르세요. 모든 패턴은 위의 리더/라이터 설정을 가정하고 제한 규칙을 따릅니다.

패턴: 메시지 필터링

속성/키 기준(ser/de 없음 — 권장). RouteOnAttribute를 사용합니다. FlowFile 속성만 읽으므로 콘텐츠 파싱 비용이 없습니다.

  1. RouteOnAttribute 프로세서를 추가합니다.

  2. 값이 Expression Language 조건인 동적 속성을 추가합니다. 각각이 관계(relationship)가 됩니다.

  3. 원하는 관계를 다음 단계에 연결합니다. unmatched(버리려면)를 auto-terminate에, 또는 Output Port에(유지하려면) 연결합니다.

콘텐츠 기준(ser/de 필요). SQL WHERE 절로 QueryRecord를 사용합니다.

Property Value
Record Reader The existing JsonTreeReader.
Record Writer The existing JsonRecordSetWriter.
filtered (dynamic) SELECT * FROM FLOWFILE WHERE

filtered 관계를 다음 단계로 라우팅하고 original을 auto-terminate하세요.

패턴: 매핑 / 필드 변환

평탄화/이름 변경/제거/기본값 연산에는 JoltTransformRecord(레코드 인식)를 사용합니다.

Property Value
Record Reader The existing JsonTreeReader.
Record Writer The existing JsonRecordSetWriter.
Jolt Transform jolt-transform-chain
Jolt Specification A Chain spec (see below).

연산을 하나의 Chain 사양에 결합해 단일 ser/de 패스로 유지하세요:

[
  {"operation": "default", "spec": {"fieldName": "defaultValue"}},
  {"operation": "shift",   "spec": {"oldName": "newName", "*": "&"}},
  {"operation": "remove",  "spec": {"unwantedField": ""}}
]

Note

JoltTransformJSON이 아니라 JoltTransformRecord를 사용하세요. JoltTransformRecord는 레코드 인식이며 구성된 RecordReader/RecordWriter로 NDJSON을 올바르게 처리합니다. JoltTransformJSON을 사용하지 마세요 — FlowFile 전체를 단일 JSON 문서로 취급해 NDJSON 입력에서 실패합니다.

더 단순한 연산(필드 추가/이름 변경/제거, 기본값 설정)에는 Jolt 사양 문법 대신 RecordPath 표현식의 UpdateRecord가 간단한 대안입니다. 콘텐츠 기반 필터링에는 QueryRecord(SQL WHERE)가 대안입니다.

패턴: 토픽-테이블 매핑

Kafka 토픽(Kinesis면 소스 스트림)을 기준으로 메시지를 다른 Snowflake 테이블로 라우팅합니다. ConsumeKafka는 kafka.topic 속성을, ConsumeKinesis는 aws.kinesis.stream.name을 자동으로 설정하므로 속성 전용입니다(ser/de 없음). 아래 단계는 kafka.topic을 사용합니다 — Kinesis에서는 aws.kinesis.stream.name으로 대체하세요.

  1. 동적 속성이 있는 UpdateAttribute 프로세서를 추가합니다:
Property Value
table.name Kafka: ${kafka.topic:replaceByPattern(#{'Topic To Table Map'})}Kinesis: ${aws.kinesis.stream.name:replaceByPattern(#{'Topic To Table Map'})}
  1. 커넥터의 파라미터 컨텍스트에 Topic To Table Map 파라미터를 추가합니다. 형식: 쉼표로 구분된 topic:table 쌍. 테이블 이름은 유효한 따옴표 없는 Snowflake 식별자여야 합니다. 정규식 패턴은 토픽을 단일 테이블에 매핑해야 합니다. 비어 있거나 일치하지 않으면 토픽 이름이 테이블 이름으로 사용됩니다.
  • 명시적: topic1:low_range,topic2:low_range,topic5:high_range

  • 정규식: topic[0-4]:low_range,topic[5-9]:high_range

  1. PublishSnowpipeStreaming을 업데이트합니다:
Property Value
Table ${table.name}

패턴: null/빈 필드의 기본값

default 연산의 JoltTransformRecord(가능하면 위 Chain 사양에 결합) 또는 RecordPath의 UpdateRecord를 사용합니다:

Property Value
Record Reader The existing JsonTreeReader.
Record Writer The existing JsonRecordSetWriter.
Replacement Value Strategy record-path-value
/field_name (dynamic) replaceNull(/field_name, 'default_value')

패턴: 여러 테이블로의 콘텐츠 기반 라우팅

메시지 콘텐츠의 필드 값에 따라 레코드를 다른 테이블로 라우팅합니다.

단일 FlowFile이 서로 다른 라우팅 필드 값을 가진 레코드를 포함할 수 있으므로 PartitionRecord는 항상 필요합니다 — 지름길이 없습니다.

1단계 — PartitionRecord 라우팅 필드로 분할하고 값을 FlowFile 속성으로 설정:

Property Value
Record Reader The existing JsonTreeReader.
Record Writer The existing JsonRecordSetWriter.
(dynamic) /routing_field

2단계(선택) — UpdateAttribute 다른 경우에만 속성을 테이블 이름으로 매핑:

Property Value
table.name ${}

속성 값 자체가 테이블 이름이라면 이것을 건너뛰고 PublishSnowpipeStreaming을 속성에 직접 가리키세요.

3단계(선택) — RouteOnAttribute 원하지 않는 값을 버려야 할 때만:

Property Value
matched (dynamic) ${table.name:isEmpty():not()}

표현식을 만족하는 레코드는 matched 관계로 흐릅니다. 이를 Output Port에 연결하세요. 나머지는 내장 unmatched 관계로 빠지며, 자동 종료해 그 레코드들을 버립니다.

4단계 — PublishSnowpipeStreaming: Table = ${table.name}(또는 라우팅 속성 직접)로 설정. PublishSnowpipeStreaming은 속성이 해석하는 테이블에 각 FlowFile을 씁니다.

Note

데이터베이스·스키마도 동적으로. PublishSnowpipeStreaming은 Table뿐 아니라 Database·Schema·Pipe 속성에 Expression Language(FlowFile 속성)를 지원합니다. 예를 들어 Database = ${target.db}, Schema = ${target.schema}, Table = ${table.name}(각 속성은 업스트림의 UpdateAttribute / PartitionRecord가 설정)으로 완전히 동적인 대상을 라우팅할 수 있습니다. 각각 FlowFile별로 평가되므로 단일 프로세서가 데이터베이스·스키마·테이블을 가로질러 팬아웃할 수 있습니다.

Need Pipeline
Route to multiple tables (all values valid) PartitionRecord -> Output Port, with PublishSnowpipeStreaming Table = ${field}
Route to multiple tables + rename attribute PartitionRecord -> UpdateAttribute -> Output Port
Route to multiple tables + drop some values PartitionRecord -> UpdateAttribute -> RouteOnAttribute -> Output Port

모든 연결에 FIFO를 설정하세요.

패턴: 사용자 정의 Groovy 스크립트(모든 상황)

위 패턴에 맞지 않는 모든 것에는 ExecuteGroovyScript를 사용하세요. 모든 제한 규칙이 여전히 적용됩니다: 연결에 FIFO, 출력 FlowFile당 하나의 테이블, 순서 보존.

Warning

원래 속성 보존. 스크립트는 들어오는 속성을 제거하거나 덮어써서는 안 됩니다. Kafka에는 kafka.topic과 kafka.partition이, Kinesis에는 aws.kinesis.stream.name과 aws.kinesis.shard.id가 포함됩니다. 이들은 PublishSnowpipeStreaming 채널 이름에 사용되며 보존되어야 합니다 — 커넥터는 at-least-once 전달을 제공하고 오프셋을 추적하지 않습니다. 출력 FlowFile을 배출할 때 항상 들어오는 FlowFile에서 속성을 상속하세요.

흐름을 활성화하기 전에 스크립트를 엣지 케이스(null, 누락 필드, 타입 불일치)에 대해 검증하세요.

변환 결합

그룹 안에서 여러 프로세서를 체인으로 연결합니다:

Input Port -> Processor A -> Processor B -> ... -> Output Port
  • 순서: 속성 전용 프로세서(ser/de 없음)를 콘텐츠 인식 프로세서 앞에 배치해 원하지 않는 데이터를 일찍 버리고, 값비싼 ser/de를 실제로 기록될 레코드에만 실행하세요.

  • ser/de 최소화: 가능하면 콘텐츠 연산을 단일 프로세서로 결합(하나의 QueryRecord 또는 체인된 JoltTransformRecord). PartitionRecord는 FlowFile 경계를 바꾸므로 분리해 두세요.

  • 연결: 모두 FIFO. 실패를 auto-terminate하거나 dead-letter 출력으로 라우팅.

파라미터화

변환을 연결한 뒤 적절한 하드코딩 값을 커넥터의 파라미터 컨텍스트로 옮겨 프로세서를 편집하지 않고 흐름을 재구성할 수 있게 하세요.

좋은 후보: 연결 문자열/URL, 토픽-테이블 매핑, 자격 증명, 임계값·상수.

적합하지 않음(인라인 유지): Groovy 스크립트, Jolt 사양, 필터/라우팅 조건, 큰 스키마 정의.

검증

모든 프로세서를 만들고 모든 연결을 걸고 PublishSnowpipeStreaming을 업데이트한 뒤:

  1. 먼저 모든 컨트롤러 서비스를 활성화하세요. 비활성화된 서비스를 참조하는 프로세서는 INVALID로 표시되므로 검증 전에 모든 서비스를 활성화하세요.

  2. 프로세스 그룹을 검증하고 검증 실패를 해결합니다.

  3. PublishSnowpipeStreaming이 올바른 Table 값(예: ${table.name})을 참조하는지 확인합니다.

  4. 흐름을 시작합니다.

문제 해결

Symptom Likely cause
Data loss or duplicate data FIFO prioritizer missing on a connection.
Out-of-order delivery Missing FIFO prioritizer on a connection between processors.
Processor shows INVALID A referenced controller service is disabled, or a required property is missing — enable services first, then re-validate.
Records written to the wrong table PublishSnowpipeStreaming Table not set to the routing attribute, or PartitionRecord / UpdateAttribute not setting the expected attribute.
Schema re-inferred on every message (slow) VolatileSchemaCache not configured on the JsonTreeReader (see Reader and writer setup).
Lost kafka.* attributes after a Groovy step The script created new FlowFiles without inheriting the incoming attributes.

함께 보기

더 알아보기 (Learn more)