CREATE PIPE

CREATE PIPE

Snowpipe가 수집 큐에서 데이터를 로드하거나, Snowpipe Streaming의 고성능 아키텍처가 스트리밍 소스에서 데이터를 테이블에 직접 로드하는 데 사용되는 COPY INTO <table> 문을 정의하기 위해 시스템에 새 파이프(pipe)를 만드는 명령이에요.

출처: 문서

본문

Snowpipe가 수집 큐에서 데이터를 로드하거나, Snowpipe Streaming의 고성능 아키텍처가 스트리밍 소스에서 데이터를 테이블에 직접 로드하는 데 사용되는 COPY INTO <table> 문을 정의하기 위해 시스템에 새 파이프를 만들어요.

이 명령은 다음 변형(variant)을 지원해요.

  • CREATE OR ALTER PIPE: 파이프가 없으면 만들고, 있으면 기존 파이프를 수정해요.

함께 보기: ALTER PIPE, DROP PIPE, SHOW PIPES, DESCRIBE PIPE

구문 (Syntax)

CREATE [ OR REPLACE ] PIPE [ IF NOT EXISTS ] <name>
  [ AUTO_INGEST = [ TRUE | FALSE ] ]
  [ ERROR_INTEGRATION = <integration_name> ]
  [ AWS_SNS_TOPIC = '<string>' ]
  [ INTEGRATION = '<string>' ]
  [ COMMENT = '<string_literal>' ]
  AS <copy_statement>

변형 구문 (Variant syntax)

CREATE OR ALTER PIPE

존재하지 않으면 새 파이프를 만들고, 존재하면 기존 파이프를 문에 정의된 파이프로 변환해요. CREATE OR ALTER PIPE 문은 CREATE PIPE 문의 구문 규칙을 따르고, ALTER PIPE 문과 동일한 제약을 가져요.

파이프를 수정할 때 다음 변경이 지원돼요.

  • PIPE_EXECUTION_PAUSED로 파이프 일시 중지 또는 재개
  • COMMENT 추가·업데이트·제거

자세한 내용은 CREATE OR ALTER <object>를 참고해요.

CREATE OR ALTER PIPE <name>
  [ AUTO_INGEST = [ TRUE | FALSE ] ]
  [ ERROR_INTEGRATION = <integration_name> ]
  [ AWS_SNS_TOPIC = '<string>' ]
  [ INTEGRATION = '<string>' ]
  [ PIPE_EXECUTION_PAUSED = TRUE | FALSE ]
  [ COMMENT = '<string_literal>' ]
  AS <copy_statement>

참고: <copy_statement>는 두 가지 유형의 데이터 소스와 함께 사용할 수 있어요.

  • 스테이지 위치: COPY INTO mytable FROM @mystage ...
  • 스트리밍 소스: COPY INTO mytable FROM (SELECT ... FROM TABLE(DATA_SOURCE(TYPE => 'STREAMING')))

필수 매개변수 (Required parameters)

name 파이프의 식별자로, 파이프가 만들어지는 스키마 안에서 고유해야 해요.

식별자는 반드시 알파벳 문자로 시작해야 하며, 전체 식별자 문자열이 큰따옴표로 묶이지 않는 한 공백이나 특수 문자를 포함할 수 없어요 (예: "My object"). 큰따옴표로 묶인 식별자는 대소문자를 구분해요.

자세한 내용은 식별자 요구 사항(Identifier requirements)을 참고해요.

copy_statement 큐에 있는 파일에서 데이터를 Snowflake 테이블로 로드하는 데 사용되는 COPY INTO <table> 문이에요. 이 문은 파이프의 텍스트/정의 역할을 하며 SHOW PIPES 출력에 표시돼요.

참고: 현재 Snowpipe의 copy_statement에 다음 함수를 사용하는 것은 권장하지 않아요.

  • CURRENT_DATE
  • CURRENT_TIME
  • CURRENT_TIMESTAMP
  • GETDATE
  • LOCALTIME
  • LOCALTIMESTAMP
  • SYSDATE
  • SYSTIMESTAMP

이 함수들로 삽입된 시간 값이 COPY_HISTORY 함수나 COPY_HISTORY 뷰가 반환하는 LOAD_TIME 값보다 몇 시간 더 일찍일 수 있는 알려진 문제가 있어요.

대신 더 정확한 레코드 로딩 표현을 제공하는 METADATA$START_SCAN_TIME을 쿼리하는 것을 권장해요.

선택 매개변수 (Optional parameters)

AUTO_INGEST = TRUE | FALSE 내부 또는 외부 스테이지에서 데이터 파일을 자동으로 로드할지 지정해요.

  • TRUE는 자동 데이터 로딩을 활성화해요. Snowpipe는 외부 스테이지(Amazon S3, Google Cloud Storage, Microsoft Azure)에서의 로딩을 지원해요.
  • FALSE는 자동 데이터 로딩을 비활성화해요. 데이터 파일을 로드하려면 Snowpipe REST API 엔드포인트에 호출을 해야 해요.
    • Snowpipe는 내부 스테이지(즉 Snowflake 이름 지정 스테이지 또는 테이블 스테이지, 사용자 스테이지는 아님) 또는 외부 스테이지(Amazon S3, Google Cloud Storage, Microsoft Azure)에서의 로딩을 지원해요.

ERROR_INTEGRATION = 'integration_name' Snowpipe가 클라우드 메시징 서비스에 오류 알림을 보내도록 구성할 때만 필요해요. 메시징 서비스와 통신하는 데 사용되는 알림 인티그레이션의 이름을 지정해요. 자세한 내용은 Snowpipe 오류 알림(Snowpipe error notifications)을 참고해요.

AWS_SNS_TOPIC = 'string' SNS를 사용해 Amazon S3 외부 스테이지용 AUTO_INGEST를 구성할 때만 필요해요. S3 버킷용 SNS 토픽의 ARN(Amazon Resource Name)을 지정해요. CREATE PIPE 문은 지정된 SNS 토픽에 Amazon Simple Queue Service(SQS) 큐를 구독해요. 파이프는 SNS 토픽을 통한 이벤트 알림으로 트리거되는 파일을 수집 큐에 복사해요. 자세한 내용은 Amazon S3용 Snowpipe 자동화(Automating Snowpipe for Amazon S3)를 참고해요.

INTEGRATION = 'string' Google Cloud Storage 또는 Microsoft Azure 외부 스테이지용 AUTO_INGEST를 구성할 때만 필요해요. 스토리지 큐에 접근하는 데 사용되는 기존 알림 인티그레이션을 지정해요. 자세한 내용은 다음을 참고해요.

  • Google Cloud Storage용 Snowpipe 자동화
  • Microsoft Azure Blob Storage용 Snowpipe 자동화

인티그레이션 이름은 모두 대문자로 입력해야 해요.

COMMENT = 'string_literal' 파이프에 대한 설명(comment)을 지정해요.

  • 기본값: 값 없음

고성능 아키텍처의 Snowpipe Streaming용 파이프

스테이지 파일 위치 없이 Snowpipe Streaming API에서 직접 데이터를 로드하도록 Snowpipe Streaming용 파이프를 정의할 수 있어요. 이 방법은 저지연, 행 기반 수집을 위해 설계됐어요.

스트리밍 파이프의 COPY INTO 문은 FROM 절에 TYPE => 'STREAMING' 인자를 가진 DATA_SOURCE 테이블 함수를 사용해야 해요.

참고: 스트리밍용으로 만든 파이프는 AUTO_INGEST 매개변수나 FROM @stage 절이 필요하지 않아요.

스트리밍 파이프 정의 안의 copy_statement는 API에서 받은 데이터를 변환하고 로드하는 데 사용돼요.

Snowpipe Streaming은 각 테이블에 대한 기본 파이프도 제공하며, 이는 요청 시 자동으로 만들어져요. 비행 중(in-flight) 변환이나 사전 클러스터링 같은 기능이 필요한 경우에만 사용자 지정 파이프를 만들면 돼요.

예시는 CREATE OR ALTER PIPE를 참고해요.

사용 메모 (Usage notes)

이 SQL 명령은 다음 최소 권한을 요구해요.

권한 (Privilege) 객체 (Object) 비고
CREATE PIPE Schema
USAGE 파이프 정의의 Stage 외부 스테이지만
USAGE Integration Snowpipe 오류 알림 수신에 필요
READ 파이프 정의의 Stage 내부 스테이지만
SELECT, INSERT 파이프 정의의 Table
OWNERSHIP Pipe 기존 파이프에 대해 CREATE OR ALTER PIPE 문을 실행하는 데 필요해요.

스키마 객체에 대한 SQL 작업은 그 객체를 포함하는 데이터베이스와 스키마에 대한 USAGE 또는 다른 권한도 필요해요.

다음을 제외한 모든 COPY INTO <table> 복사 옵션이 지원돼요.

  • FILES = ( 'file_name1' [ , 'file_name2', ... ] )
  • ON_ERROR = ABORT_STATEMENT
  • SIZE_LIMIT = num
  • PURGE = TRUE | FALSE (즉, 로딩 중 자동 삭제)
  • FORCE = TRUE | FALSE

(로드된 후) 내부(즉 Snowflake) 스테이지에서 파일을 REMOVE 명령으로 수동으로 제거할 수 있다는 점에 주의해요.

  • RETURN_FAILED_ONLY = TRUE | FALSE
  • VALIDATION_MODE = RETURN_n_ROWS | RETURN_ERRORS | RETURN_ALL_ERRORS

PATTERN = 'regex_pattern' 복사 옵션은 정규 표현식으로 로드할 파일 집합을 필터링해요. 패턴 일치는 AUTO_INGEST 매개변수 값에 따라 다음과 같이 동작해요.

  • AUTO_INGEST = TRUE: 정규 표현식이 COPY INTO <table> 문의 스테이지와 선택적 경로(즉 클라우드 스토리지 위치)의 파일 목록을 필터링해요.
  • AUTO_INGEST = FALSE: 정규 표현식이 Snowpipe REST API의 insertFiles 엔드포인트에 호출되어 제출된 파일 목록을 필터링해요.

Snowpipe는 스테이지 정의에서 어떤 경로 세그먼트든 스토리지 위치에서 잘라내고, 나머지 경로 세그먼트와 파일 이름에 정규 표현식을 적용해요. 스테이지 정의를 보려면 해당 스테이지에 대해 DESCRIBE STAGE 명령을 실행해요. URL 속성은 버킷 또는 컨테이너 이름과 0개 이상의 경로 세그먼트로 구성돼요. 예를 들어 COPY INTO <table> 문의 FROM 위치가 @s/path1/path2/이고 스테이지 @s의 URL 값이 s3://mybucket/path1/이라면, Snowpipe는 FROM 절의 스토리지 위치에서 s3://mybucket/path1/path2/를 잘라내고 나머지 경로의 파일 이름에 정규 표현식을 적용해요.

⚠️ 중요: Snowflake는 비용·이벤트 노이즈·지연 시간을 줄이기 위해 Snowpipe용 클라우드 이벤트 필터링을 활성화할 것을 권장해요. 클라우드 공급자의 이벤트 필터링 기능이 충분하지 않을 때만 PATTERN 옵션을 사용해요. 각 클라우드 공급자의 이벤트 필터링 구성에 대한 자세한 내용은 다음 페이지를 참고해요.

  • Amazon S3: 객체 키 이름 필터링을 사용한 이벤트 알림 구성
  • Microsoft Azure Event Grid: Event Grid 구독의 이벤트 필터링 이해
  • Google Cloud Pub/Sub: 메시지 필터링

컬럼 재정렬·컬럼 생략·캐스트(즉 로딩 중 데이터 변환)를 위한 COPY 문의 소스로 쿼리를 사용하는 것은 지원돼요. 사용 예시는 로딩 중 데이터 변환(Transform data during a load)을 참고해요. 단순 SELECT 문만 지원된다는 점에 주의해요. WHERE 절을 사용한 필터링은 지원되지 않아요.

파이프 정의는 동적이지 않아요 (즉, 기본 스테이지나 테이블이 변경되면[이름 바꾸기, 삭제 등] 파이프가 자동으로 업데이트되지 않아요). 대신 새 파이프를 만들고 향후 Snowpipe REST API 호출에서 이 파이프 이름을 제출해야 해요.

메타데이터에 관해서는 다음 사항에 주의해요.

⚠️ 주의: 고객은 Snowflake 서비스를 사용할 때 (User 객체를 제외하고) 개인 데이터·민감 데이터·수출 통제 데이터·기타 규제 데이터를 메타데이터로 입력하지 않도록 해야 해요. 자세한 내용은 Snowflake의 메타데이터 필드를 참고해요.

OR REPLACEIF NOT EXISTS 절은 서로 배타적이에요. 같은 문에서 둘 다 사용할 수 없어요.

CREATE OR REPLACE <object> 문은 원자적(atomic)으로 동작해요. 즉, 객체를 교체할 때 기존 객체는 삭제되고 새 객체는 단일 트랜잭션 안에서 생성돼요.

⚠️ 중요: 파이프를 다시 만들면(CREATE OR REPLACE PIPE 구문), 관련 고려 사항과 모범 사례는 파이프 다시 만들기(Recreating pipes)를 참고해요.

CREATE OR ALTER PIPE

ALTER PIPE 명령의 모든 제약이 적용돼요.

기존 파이프의 COPY INTO 문(copy_statement)은 변경할 수 없어요. 파이프 정의를 변경해야 한다면 기존 파이프를 삭제하고 새로 만들어요. 관련 고려 사항은 파이프 다시 만들기를 참고해요.

기존 파이프의 AUTO_INGEST, AWS_SNS_TOPIC, INTEGRATION 속성은 변경할 수 없어요.

태그 설정·해제는 지원되지 않아요.

예시 (Examples)

현재 스키마에 mystage 스테이지에 스테이징된 파일의 모든 데이터를 mytable로 로드하는 파이프를 만들어요.

CREATE PIPE mypipe
  AS
  COPY INTO mytable
  FROM @mystage
  FILE_FORMAT = (TYPE = 'JSON');

이전 예시와 같지만 데이터 변환이 있어요. 스테이지된 파일의 4번째와 5번째 컬럼의 데이터만 역순으로 로드해요.

CREATE PIPE mypipe2
  AS
  COPY INTO mytable(C1, C2)
  FROM (SELECT $5, $4 FROM @mystage)
  FILE_FORMAT = (TYPE = 'JSON');

데이터에 표시된 해당 컬럼과 일치하는 대상 테이블의 컬럼에 모든 데이터를 로드하는 파이프를 만들어요. 컬럼 이름은 대소문자를 구분하지 않아요.

또한 METADATA$START_SCAN_TIME과 METADATA$FILENAME 메타데이터 컬럼에서 c1c2라는 컬럼으로 메타데이터를 로드해요.

CREATE PIPE mypipe3
  AS
  (COPY INTO mytable
    FROM @mystage
    MATCH_BY_COLUMN_NAME=CASE_INSENSITIVE
    INCLUDE_METADATA = (c1= METADATA$START_SCAN_TIME, c2=METADATA$FILENAME)
    FILE_FORMAT = (TYPE = 'JSON'));

메시징 서비스에서 받은 이벤트 알림을 사용해 데이터를 자동 로드하기 위한 파이프를 현재 스키마에 만들어요.

Amazon S3:

CREATE PIPE mypipe_s3
  AUTO_INGEST = TRUE
  AWS_SNS_TOPIC = 'arn:aws:sns:us-west-2:001234567890:s3_mybucket'
  AS
  COPY INTO snowpipe_db.public.mytable
  FROM @snowpipe_db.public.mystage
  FILE_FORMAT = (TYPE = 'JSON');

Google Cloud Storage:

CREATE PIPE mypipe_gcs
  AUTO_INGEST = TRUE
  INTEGRATION = 'MYINT'
  AS
  COPY INTO snowpipe_db.public.mytable
  FROM @snowpipe_db.public.mystage
  FILE_FORMAT = (TYPE = 'JSON');

Microsoft Azure:

CREATE PIPE mypipe_azure
  AUTO_INGEST = TRUE
  INTEGRATION = 'MYINT'
  AS
  COPY INTO snowpipe_db.public.mytable
  FROM @snowpipe_db.public.mystage
  FILE_FORMAT = (TYPE = 'JSON');

내부 이름 지정 스테이지:

mystage라는 내부 이름 지정 스테이지의 모든 데이터 파일을 자동으로 로드하는 파이프를 현재 스키마에 만들어요.

CREATE PIPE mypipe_aws
  AUTO_INGEST = TRUE
  AS
  COPY INTO snowpipe_db.public.mytable
  FROM @snowpipe_db.public.mystage
  FILE_FORMAT = (TYPE = 'JSON');

고성능 아키텍처의 Snowpipe Streaming:

기본 스트리밍 파이프를 만들어요.

CREATE OR REPLACE PIPE my_streaming_pipe
AS COPY INTO my_table
  FROM (SELECT $1, $1:c1, $1:ts FROM TABLE(DATA_SOURCE(TYPE => 'STREAMING')));

SELECT 절에 컬럼 표현식을 지정해 비행 중 변환이 있는 스트리밍 파이프를 만들어요.

CREATE OR REPLACE PIPE my_pipe_with_transforms
AS COPY INTO my_table (col1, col2, col3)
  FROM (
    SELECT
      $1:field1::STRING AS col1,
      $1:field2::NUMBER AS col2,
      CURRENT_TIMESTAMP() AS col3
    FROM TABLE (DATA_SOURCE(TYPE => 'STREAMING'))
  );

향상된 쿼리 성능을 위해 사전 클러스터링이 활성화된 스트리밍 파이프를 만들어요. 대상 테이블에 클러스터링 키가 정의되어 있어야 해요.

CREATE OR REPLACE PIPE my_pipe_with_clustering
AS COPY INTO my_table
  FROM (SELECT $1, $1:c1, $1:ts FROM TABLE(DATA_SOURCE(TYPE => 'STREAMING')))
  CLUSTER_AT_INGEST_TIME = TRUE;

Apache Iceberg v3 지원:

참고: Snowflake가 지원하는 다른 Iceberg v3 기능에 대한 자세한 내용은 Apache Iceberg™ 테이블: Apache Iceberg™ v3 지원을 참고해요.

다음 예시는 Snowflake 관리 테이블과 외부 관리 테이블 모두에 대해 Iceberg v3 테이블의 파일에서 데이터를 로드해요.

CREATE PIPE mypipe
  AUTO_INGEST = TRUE
  INTEGRATION = 'MYINT'
  AS
  COPY INTO snowpipe_db.public.my_v3_iceberg_table
  FROM @snowpipe_db.public.mystage
  FILE_FORMAT = (TYPE = 'JSON');

CREATE OR ALTER PIPE

새 파이프를 만들거나 기존 파이프를 일시 중지해요.

CREATE OR ALTER PIPE mypipe
  PIPE_EXECUTION_PAUSED = TRUE
  COMMENT = 'Loads data from mystage into mytable'
  AS
  COPY INTO mytable
  FROM @mystage
  FILE_FORMAT = (TYPE = 'JSON');

더 알아보기 (Learn more)