Elastic Channels 모범 사례
Elastic Channels 모범 사례
Elastic Channels를 가장 효율적이고 안전하게 사용하기 위한 모범 사례를 알려드릴게요. SDK 자동 일괄 처리 활용, 승인 처리, 중복 방지, 장애 대비 이벤트 보관 등을 다룹니다.
출처: Snowflake 문서
본문
SDK가 자동으로 일괄 처리하게 하세요
행이 도착하면 바로 append하세요. Java, Python, Node.js SDK는 시간·크기 임계값을 사용해 내부적으로 append를 버퍼링·결합합니다. SDK가 압축과 Snowflake 전송도 처리해요. SDK에 제출하기 위해 행을 모아둘 때까지 기다리지 마세요.
개별 행을 제출할 때는
미완료 승인 수를 제한하세요
각 append가 승인되기를 기다리지 않고 다음 행을 제출하세요. 미완료 승인 수와 미승인 이벤트에 대해 보관하는 바이트를 제한하세요. 미완료 작업 한도나 경과 시간을 기준으로 주기적으로 모든 보류 중인 append가 내구성 있게 승인될 때까지 기다리세요. 수입이 일시 중지되거나 끝날 때도 기다리세요. 이는 이미 제출된 이벤트를 기다리는 것이지, 보낼 행을 모으는 것이 아니에요.
모든 관련 이벤트가 내구성 있게 승인된 후에만 애플리케이션 체크포인트를 전진시키세요. Elastic Channels는 승인 순서를 보장하지 않으므로, 이후 append의 완료만으로는 이전 append의 성공을 확립하지 못해요. 체크포인트와 플러시 동작은 Elastic Channels 작업을 참고하세요.
SDK의 버퍼가 가득 차서 새 append를 거부하면 수입을 일시 중지하고 거부된 이벤트를 보관하세요. 보류 중인 append가 완료되게 둔 다음, 거부된 append를 재시도하기 전에 잠시 기다리세요.
내구성 승인에 시간을 허용하세요
SDK는 일시적인 네트워크·서비스 실패를 자동으로 재시도하므로, 재시도나 서비스 지연 중에는 승인에 더 오래 걸릴 수 있어요. 느린 승인을 실패한 append로 취급하는 10초 같은 공격적인 per-append 타임아웃은 피하세요. 애플리케이션이 제한된 대기를 필요로 한다면 원래 Future나 Promise를 보관하고 그 결과를 계속 추적하세요. 미완료 append와 보관된 바이트를 제한하고, 그 한도에 도달하면 보류 중인 데이터를 재제출하는 대신 수입을 일시 중지하세요. 언어별 지침은 호출자 대기 타임아웃을 참고하세요.
안정적인 이벤트 식별자를 사용하세요
Elastic 전달은 at-least-once입니다. 프로듀서가 모호한 실패(타임아웃, 응답 없음, 5xx) 후 재시도하면 Snowflake가 원래 요청을 이미 수락했을 수 있고 대상 테이블에 중복 행이 있을 수 있어요.
각 행에 안정적인 이벤트 식별자(예: 소스에서 생성한 UUID, 또는 소스 시스템·파티션·offset의 조합)를 포함하세요. 중복이 문제가 될 때 이 식별자를 다운스트림 중복 제거에 사용하세요.
-- 안정적인 event_id에 대한 예시 다운스트림 중복 제거
SELECT DISTINCT event_id, * FROM MY_TABLE;
Append token은 콜백에서 반환되어 제출된 메시지에 결과를 매칭하는 식별자입니다. Snowflake로 전송되지 않으며 중복을 방지하지 않아요.
미승인 데이터를 보호하세요
Snowflake가 append를 승인하면 프로듀서는 보관한 복사본을 해제할 수 있어요. 그때까지는 이벤트를 계속 유지하세요. 장애 중에는 수입을 일시 중지하고 보류 중인 이벤트를 보관하세요. 애플리케이션이 계속 수집해야 한다면, SDK에 제출하기 전에 들어오는 이벤트를 내구성 있는 스토리지에 보관하세요.
주 경로는 계속 직접적입니다: 프로듀서 → Snowpipe Streaming SDK → Snowflake. SDK는 프로세스 메모리에 append를 버퍼링하고 일시적 실패를 재시도하지만, 복구용 내구성 로컬 스토리지는 제공하지 않아요. 전달 요구사항에 따라 보존을 선택하세요.
| 프로듀서 요구사항 | 권장 패턴 | Kafka 필요? |
|---|---|---|
| 수입을 일시 중지하거나 이벤트를 재생성할 수 있음 | 직접 SDK 수입. 장애 중 수입을 일시 중지하고 이벤트가 내구성 있게 승인될 때까지 보관. 복구가 프로듀서 재시작을 견뎌야 하면 소스 재생 사용. | 아니요 |
| 장애 중에도 이벤트 수락을 계속해야 함 | 보류 중인 이벤트용 애플리케이션 관리 내구성 스토리지가 있는 직접 SDK 수입. | 아니요 |
| 이미 공유 메시징 시스템이 필요함 | 여러 소비자로의 전달, 재생, 보존을 위해 Kafka를 유지하고 Snowflake에 연결. | 해당 요구사항에 적합. 수입에 필수는 아님 |
수입을 일시 중지하면 추가 작업이 제한되지만, 메모리 전용 이벤트를 크래시에서 살아남게 하지는 않아요. 소스가 재생 가능하다면 레코드를 버리거나 미승인 이벤트를 지나 체크포인트를 전진시키지 마세요.
이미 여러 소비자, 공유 재생, 보존을 위해 Kafka를 사용한다면 유지하고 그 복사본을 Snowflake로 스트리밍하세요. 단순히 이벤트를 테이블로 스트리밍하기 위해 Kafka를 새로 도입할 필요는 없어요.
장애 중에도 수집을 계속하세요
재생할 수 없는 소스에서 손실에 민감한 수집을 하려면, 그 이벤트에 대한 책임을 받아들이기 전에 내구성 있는 스토리지에 이벤트를 영속화하세요. SDK 버퍼가 가득 찰 때까지 기다렸다가 영속화를 시작하지 마세요. 그 시점 전에 크래시가 나면 수락된 이벤트를 잃을 수 있어요. 이 프로듀서 측 보존은 별도의 브로커나 커넥터 플릿 없이 Snowflake로의 직접 네트워크 경로를 보존할 수 있어요.
내구성 승인 후에만 영속화된 레코드를 제거하세요. 재시작 후에는 안정적인 이벤트 ID를 사용해 미해결 레코드를 재생하세요. 승인이 프로듀서가 기록하기 전에 유실됐을 수도 있기 때문이에요.
로컬 스토리지를 제한하고 가득 찼을 때 무슨 일이 일어날지 정의하세요. 로컬 디스크는 호스트 손실을 보호하지 못하며, 일시적인 컨테이너 파일시스템은 재시작을 견디지 못할 수 있어요. 보존이 호스트 손실에서 살아남아야 한다면 필요한 복제·내구성 보장이 있는 스토리지를 사용하세요. 일시 중지도, 재생도, 영속화도 할 수 없는 프로듀서는 장애 중 무손실 수집을 보장할 수 없어요.
콜백을 짧게 유지하세요
콜백으로 추적되는 첫 append 전에
콜백은 다음을 해야 합니다:
- 빠르게 반환할 것
- 저렴한 부기(bookkeeping)만 할 것 (예: 로컬 카운터에 승인된 token 기록)
- 차단 I/O, 네트워크 호출, 재시도, 조정은 애플리케이션 관리 큐나 실행기로 넘길 것
느린 콜백은 채널의 모든 append의 승인을 지연시킵니다. 던져진 콜백 예외는 포착되어 기록되므로, 하나의 잘못된 핸들러가 승인 경로를 멈추게 하지는 않아요.
REST로 처리량을 위한 행 일괄 처리
직접 REST 클라이언트는 행을 스스로 그룹화하고 보내야 합니다. 행을 NDJSON(newline-delimited JSON, 한 줄에 JSON 객체 하나)을 사용해 요청으로 결합하고 ZSTD나 Gzip 압축으로 요청 오버헤드를 줄이세요. 요청 크기를 제한하고 경과 시간에 따라 부분 배치를 플러시해서 저용량 소스가 전체 배치를 무한정 기다리지 않게 하세요.
Elastic REST 요청에는 4 MB 페이로드 제한이 있습니다(압축을 쓰면 압축 후 네트워크로 보내지는 페이로드 크기). ZSTD나 Gzip 압축을 사용해 페이로드 크기를 줄이세요. 이는 REST 요청 제한이지 애플리케이션 측 SDK 배치의 목표가 아니에요.
페이로드가 이미 일치하는 형식으로 압축되어 있을 때만
- ZSTD의 경우
Content-Encoding: zstd - Gzip의 경우
Content-Encoding: gzip
안전하게 재시도하기
직접 REST 클라이언트는 재시도 지연에 무작위 변동(jitter)이 있는 지수 백오프(exponential backoff)를 사용해 일시적 오류(HTTP 429, 500, 503)를 재시도해야 합니다. 네트워크 타임아웃이나 응답 없음은 모호합니다. 원래 요청이 이미 수락됐을 수 있어요.
REST 요청의 경우, 같은 rowset(행 배치)의 재시도마다 같은
SDK append의 경우 원래 append가 보류 중인 동안 SDK가 일시적 실패를 재시도하게 두세요. 호출자가 부과한 타임아웃은 종료적인 append 실패가 아닙니다. SDK가 종료적 실패를 보고하면 재시도하기 전에 분류하세요. 검증·권한 부여 문제를 고치고, 무효화된 클라이언트를 다시 만드세요. 모호한 결과에는 안정적인 이벤트 ID를 사용해 미해결 이벤트만 복구하세요. Elastic Channels 오류 처리를 참고하세요.
정상 종료(Graceful shutdown)
프로듀서를 멈추기 전에:
- 새 소스 작업 수락을 중지하세요.
- 버퍼링된 데이터를 밀어내려면
initiateFlush() 를 호출하세요. - 플러시가 완료될 때까지 기다리세요(
waitForFlush 타임아웃 포함). - 클라이언트를 닫으세요.
REST 프로듀서의 경우 프로세스를 멈추기 전에 모든 in-flight 요청이 성공적인 HTTP 응답을 받았는지 확인하세요. 코드 예제는 Elastic Channels 작업을 참고하세요.
효율적인 수입을 위해 MATCH_BY_COLUMN_NAME 사용
스트리밍 pipe를
반구조적 데이터에 네이티브 데이터 타입 사용
반구조적 데이터를 직렬화된 JSON 문자열 대신 네이티브 언어 객체(Java
Java:
// 권장: SDK가 Map을 구조적 VARIANT로 변환
row.put("payload", Map.of("event_id", 101, "status", "active"));
Python:
# 권장: SDK가 dict를 구조적 VARIANT로 변환
row["payload"] = {"event_id": 101, "status": "active"}
Node.js:
// 권장: SDK가 객체를 구조적 VARIANT로 변환
const row = { payload: { event_id: 101, status: "active" } };
Java 배포를 위한 JVM 힙 크기 구성
Java SDK에는 JVM 힙 밖에서 메모리를 할당하는 네이티브 Rust 컴포넌트가 포함되어 있습니다. JVM 힙을 사용 가능한 메모리의 약 50%로 제한하세요. 8 GB RAM 호스트의 경우
MAVEN_OPTS="-Xmx4g" mvn exec:java -Dexec.mainClass="com.example.Main"
# 또는
java -Xmx4g -jar your-app.jar