Named Channel과 클래식 SDK 비교

Named Channel과 클래식 SDK 비교

이 섹션에서는 Snowpipe Streaming Classic SDK와 고성능 SDK의 Named Channel API 사이의 주요 차이점을 요약합니다.

출처: Snowflake 문서

본문

클라이언트 및 채널 관리

  • OpenClient: 고성능 SDK는 DB, SCHEMA, PIPE를 지정해야 합니다. 클래식 SDK에서는 클라이언트 NAME만 지정하면 됩니다.
  • OpenChannel: 고성능 SDK는 채널 이름만 요구해 이를 단순화합니다. 클래식 SDK는 DB, SCHEMA, TABLE, ERROR_OPTION을 지정해야 합니다. 새 SDK는 또한 채널 엔티티와 상태를 포함하는 OpenChannelResult를 반환하므로, 마지막 커밋된 offset token을 얻기 위한 별도의 RPC 호출이 필요하지 않습니다.
  • offsetToken 지원: 새 openChannel 메서드는 이제 선택적 offsetToken 파라미터를 가지므로 특정 위치에서 채널을 열 수 있습니다. openChannel(String channelName, (optional) String offsetToken).

데이터 수집

  • InsertRows 이름 변경: InsertRows 메서드는 고성능 SDK에서 이제 AppendRows라고 합니다.
  • AppendResult 제거: appendRow 및 appendRows 메서드는 더 이상 AppendResult를 반환하지 않습니다. 시그니처가 void appendRow(Map row, String offsetToken) 및 void appendRows(Iterable> row, String startOffsetToken, String endOffsetToken)으로 변경되었습니다.

새로운 비동기 및 유틸리티 메서드

  • GetChannelStatus: Channel 객체에서 사용할 수 있는 새로운 API입니다.
  • waitForFlush: 클라이언트와 채널 객체 양쪽에 새로운 waitForFlush 메서드가 추가되었습니다.

Client: CompletableFuture close(boolean waitForFlush, Duration timeoutDuration)

  • Channel and Client: CompletableFuture waitForFlush((optional) Duration timeoutDuration)

  • waitForCommit: CompletableFuture<Void> waitForCommit(Predicate tokenChecker, Duration timeoutDuration)는 제공된 predicate가 성공할 때까지 최신 커밋된 offset을 폴링합니다. 이것은 Named Channel 소스 체크포인트이지 flush 작업이 아닙니다. 시간 초과나 실패 시 Future는 예외로 완료됩니다.

  • initiateFlush: 새 메서드 void initiateFlush()는 채널 또는 클라이언트에서 flush를 비동기적으로 호출합니다. 이 메서드를 사용하면 시간 초과나 크기 한도를 기다리지 않고 데이터를 flush할 수 있습니다.

데이터 타입 및 파싱

고성능 아키텍처는 ARRAY와 VARIANT 열에 네이티브 객체를 요구하며 문자열 리터럴을 자동 파싱하지 않습니다.

열 타입 Classic High-performance
OBJECT JSON 문자열을 자동 파싱합니다. 변경 없음. JSON 문자열을 자동 파싱합니다.
ARRAY 문자열을 암시적으로 파싱합니다. 예: “[1,2]”는 [1,2]가 됩니다. 타입 엄격(Type-strict). 문자열을 리터럴로 취급합니다. 예: “[1,2]”는 [“[1,2]”]가 됩니다.
VARIANT 문자열을 암시적으로 파싱합니다. 예: “true”는 true가 됩니다. 타입 엄격(Type-strict). 문자열을 리터럴로 취급합니다. 예: “true”는 “true”가 됩니다.

고성능 아키텍처에서 반정형 데이터가 올바르게 저장되도록 하려면 직렬화된 JSON 문자열 대신 네이티브 언어 객체(예: Java List/Map, Python list/dict, JavaScript Array/Object)를 전달하세요.

기타 변경 사항

  • GetLatestCommittedOffsetTokens: 이 API가 개선되었습니다. 고성능 SDK에서는 클라이언트가 열지 않은 채널에 대해서도 offset token을 가져올 수 있고 부분 실패를 허용합니다.
  • isValid 제거: isValid 메서드는 고성능 SDK에서 제거되었습니다.
  • 스키마 진화(schema evolution) 지원: 고성능 SDK는 변화하는 데이터 스키마를 자동으로 처리하는 핵심 기능인 스키마 진화 를 지원합니다.

다음 테이블은 클래식 SDK에서 고성능 SDK로의 API 변경 사항을 보여줍니다.

SnowflakeStreamingIngestClientFactory 및 SnowflakeStreamingIngestClientFactory.Builder

Classic High-performance Notes
builder(String name) builder(String clientName, String dbName, String schemaName, String pipeName) 클래식 버전의 name = 고성능 버전의 clientName.
N/A setExecutorService(ExecutorService executorService) 새 메서드. SDK가 백그라운드 작업에 사용할 ExecutorService를 지정할 수 있습니다.

SnowflakeStreamingIngestClient

Classic High-performance Notes
String getName() String getClientName() API 이름만 변경됨; 동일한 정보를 반환합니다.
N/A String getDBName() 새 API.
N/A String getPipeName() 새 API.
N/A String getSchemaName() 새 API.
SnowflakeStreamingIngestChannel openChannel(OpenChannelRequest request) OpenChannelResult openChannel(String channelName, (optional) String offsetToken) 다른 요청 인자와 반환 값.
Map getLatestCommittedOffsetTokens (List channels) Map getLatestCommittedOffsetTokens (List channelNames) 다른 요청 인자. 고성능 SDK는 다른 클라이언트가 열었고 클라이언트에 속하지 않을 수 있는 채널의 상태를 가져올 수 있습니다.
N/A ChannelStatusBatch getChannelStatus(List channelNames) 새 API.
Void dropChannel(DropChannelRequest request) Void dropChannel(String channelName) 다른 요청 인자.
Void setRefreshToken(String refreshToken) N/A 제거됨.
N/A CompletableFuture close(boolean waitForFlush, Duration timeoutDuration) 종료 프로세스를 더 세밀하게 제어하는 새 클라이언트 close 메서드. waitForFlush: 종료 전에 클라이언트가 모든 채널의 flush를 기다려야 하는지 나타내는 Boolean 파라미터. timeoutDuration: 강제 종료 전에 클라이언트가 flush 완료를 기다릴 시간을 지정하는 Duration.
N/A CompletableFuture waitForFlush((optional) Duration timeoutDuration) flush 완료를 기다리는 새 메서드. timeoutDuration: 클라이언트가 시간 초과되기 전에 기다릴 시간을 지정.
N/A void initiateFlush() 클라이언트가 flush를 비동기적으로 트리거하고 즉시 반환하는 새 메서드.

SnowflakeStreamingIngestChannel

Classic High-performance Notes
getLatestCommittedOffsetToken getLatestCommittedOffsetToken 이 API가 개선되었습니다. 고성능 SDK에서는 클라이언트가 열지 않은 채널에 대해서도 offset token을 가져올 수 있고 부분 실패를 허용합니다.
isValid N/A 제거됨.
N/A String getDBName() 새 API.
N/A String getSchemaName() 새 API.
N/A String getPipeName() 새 API.
N/A String getFullyQualifiedPipeName() 새 API.
InsertValidationResponse insertRow(Map row, String offsetToken) void appendRow(Map row, @Nullable String offsetToken) API 이름 변경. 클라이언트에서 더 이상 검증이 없으므로 응답 타입이 변경됨.
InsertValidationResponse insertRow(Iterable> row, @Nullable String startOffsetToken, @Nullable String endOffsetToken) void appendRows(Iterable> row, String startOffsetToken, String endOffsetToken) API 이름 변경. 클라이언트에서 더 이상 검증이 없으므로 응답 타입이 변경됨.
InsertValidationResponse insertRow(Iterable> row, String offsetToken) N/A 제거됨.
String getTableName() N/A 제거됨.
String getFullyQualifiedTableName() N/A 제거됨.
N/A String getPipeName() 새 API.
N/A String getFullyQualifiedPipeName() 새 API.
String getName() String getChannelName() API 이름 변경.
String getFullyQualifiedName() String getFullyQualifiedChannelName() API 이름 변경.
Map getTableSchema() N/A 제거됨.
N/A ChannelStatus getChannelStatus() 새 API.
CompletableFuture close() Void close() 반환 타입은 변경되었지만 동작은 동일합니다.
CompletableFuture close(boolean drop) Void close(boolean waitForFlush, Duration timeoutDuration) API 이름은 변경되었지만 동작은 동일합니다.
Boolean isValid() N/A 제거됨.
N/A CompletableFuture waitForFlush((optional)Duration timeoutDuration) flush 완료를 기다리는 새 메서드. timeoutDuration: 채널이 시간 초과되기 전에 기다릴 시간을 지정.
N/A CompletableFuture waitForCommit(Predicate tokenChecker, Duration timeoutDuration) 제공된 predicate가 성공할 때까지 최신 커밋된 offset을 폴링합니다. predicate는 대상 offset 너머의 커밋된 진행 상황도 처리해야 합니다. 소스 체크포인트에서 사용하고 모든 append 이후에는 사용하지 마세요. flush를 트리거하거나 SDK 배치를 제어하지 않습니다. 시간 초과나 실패 시 Future는 예외로 완료됩니다.
N/A void initiateFlush() 채널이 flush를 비동기적으로 트리거하는 새 메서드.

더 알아보기 (Learn more)