Snowpipe REST API

Snowpipe REST API

REST 엔드포인트에 호출을 함으로써 파이프와 상호작용해요. 이 주제는 수집할 파일 목록을 정의하고 로드 기록 보고서를 가져오기 위한 Snowpipe REST API를 설명해요.

Snowflake는 또한 Snowpipe REST API 작업을 단순화하는 Java 및 Python API를 제공해요.

출처: Documentation

본문

데이터 파일 수집

Snowpipe API는 수집할 파일 목록을 정의하기 위한 REST 엔드포인트를 제공해요.

엔드포인트: insertFiles

테이블에 수집될 파일에 대해 Snowflake에 알려줘요. 이 엔드포인트의 성공적인 응답은 Snowflake가 테이블에 추가할 파일 목록을 기록했음을 뜻해요. 파일이 수집되었다는 뜻은 아니에요. 자세한 내용은 아래 응답 코드를 참고해요.

대부분의 경우 Snowflake는 몇 분 안에 대상 테이블에 새 데이터를 삽입해요.

메서드: POST

POST URL:

https://{account}.snowflakecomputing.com/v1/data/pipes/{pipeName}/insertFiles?requestId={requestId}

URL 매개 변수:

  • account(필수): Snowflake 계정의 계정 식별자.
  • pipeName(필수): 대소문자를 구분하는 정규화된 파이프 이름. 예: myDatabase.mySchema.myPipe.
  • requestId(선택): 시스템을 통해 요청을 추적하는 데 사용되는 문자열. 각 요청에 UUID 같은 임의 문자열을 제공할 것을 권장해요. URL에 ?requestId=<your_uuid>처럼 추가해야 해요.

요청 헤더

  • Content-Type: — text/plain: 파일 경로와 파일 이름의 일반 텍스트 목록용(한 줄에 하나). 이 형식에서는 size 매개 변수가 허용되지 않아요. / application/json: 선택적 크기 정보가 있는 파일 목록을 포함하는 JSON 객체용.
  • Authorization: BEARER <jwt_token>

요청 본문(application/json Content-Type용)

요청 본문은 "files"라는 단일 키가 있는 JSON 객체여야 해요. 이 키와 연결된 값은 수집할 파일 각각을 나타내는 JSON 객체의 배열이에요.

{
  "files":[
    {
      "path":"filePath/file1.csv",
      "size":100
    },
    {
      "path":"filePath/file2.csv",
      "size":100
    }
   ]
}

"files" 배열의 각 요소는 다음 특성이 있는 JSON 객체예요.

  • path(필수): 스테이징된 파일의 경로와 파일 이름. 권장 모범 사례를 따라 논리적이고 세분화된 경로로 스테이지에서 데이터를 파티셔닝했다면 페이로드의 경로 값에는 스테이징된 파일의 전체 경로가 포함돼요.
  • size(선택, 성능 향상을 위해 권장): 파일의 크기(바이트 단위).

요청 본문(text/plain Content-Type용)

요청 본문은 파일 경로와 파일 이름의 일반 텍스트 목록이어야 하며, 한 줄에 하나의 항목이어야 해요.

filePath/file_a.csv
another/path/file_b.json
yet/another/file_c.txt

참고 — 게시물에는 최대 5,000개의 파일이 포함될 수 있어요. 각 파일 경로는 UTF-8로 직렬화했을 때 1024바이트 이하여야 해요.

응답 본문

응답 코드:

  • 200 — 성공. 파일이 수집할 파일 큐에 추가됨.
  • 400 — 실패. 잘못된 형식 또는 한도 초과로 인한 유효하지 않은 요청.
  • 404 — 실패. pipeName이 인식되지 않음. 이 오류 코드는 엔드포인트를 호출할 때 사용된 역할에 충분한 권한이 없는 경우에도 반환될 수 있어요. 자세한 내용은 액세스 권한 부여를 참고해요.
  • 429 — 실패. 요청 속도 한도를 초과함.
  • 500 — 실패. 내부 오류 발생.

응답 페이로드:

성공적인 API 요청(즉, 코드 200)의 경우 응답 페이로드는 JSON 형식의 requestId와 status 요소를 포함해요. 오류가 발생하면 응답 페이로드에 오류에 대한 세부 정보가 포함될 수 있어요.

{
  "requestId": "your_request_uuid",
  "status": "success"
}

파이프 정의의 COPY INTO <table> 문에 PATTERN 복사 옵션이 포함되면 unmatchedPatternFiles 특성은 헤더에 제출된 정규식과 일치하지 않아 건너뛴 파일을 나열해요.

{
  "requestId": "your_request_uuid",
  "status": "success",
  "unmatchedPatternFiles": ["some_file.txt", "another_file.dat"]
}

로드 기록 보고서

Snowpipe API는 로드 보고서를 가져오기 위한 REST 엔드포인트를 제공해요.

엔드포인트: insertReport

insertFiles를 통해 제출되었고 내용이 최근에 테이블에 수집된 파일의 보고서를 검색해요. 대용량 파일의 경우 이 보고서가 파일의 일부만일 수 있다는 점에 주의해요.

이 엔드포인트의 다음 제한 사항을 주의해요.

  • 가장 최근 10,000개의 이벤트가 보존돼요.
  • 이벤트는 최대 10분 동안 보존돼요.

이벤트는 insertFiles를 통해 제출된 파일의 데이터가 테이블에 커밋되어 쿼리에서 사용할 수 있을 때 발생해요. insertReport 엔드포인트는 UNIX 명령 tail과 같다고 생각할 수 있어요. 이 명령을 반복해서 호출하면 시간이 지남에 따라 파이프의 전체 이벤트 기록을 볼 수 있어요. 이벤트를 놓치지 않으려면 명령을 충분히 자주 호출해야 한다는 점에 주의해요. 얼마나 자주 호출해야 하는지는 insertFiles에 파일이 전송되는 비율에 따라 달라져요.

메서드: GET

GET URL:

https://<account_identifier>.snowflakecomputing.com/v1/data/pipes/<pipeName>/insertReport?requestId=<requestId>&beginMark=<beginMark>

URL 매개 변수:

  • account_identifier(필수): Snowflake의 고유 계정 식별자. 권장 형식은 *organization_name*-*account_name*. 대체 형식(리전과 클라우드 플랫폼이 있는 계정 로케이터)은 조직의 계정 이름(형식 1, 권장)을 참고해요.
  • pipeName(필수): 대소문자를 구분하는 정규화된 Snowpipe 이름. 예: myDatabase.mySchema.myPipe.
  • requestId(선택): Snowflake 시스템을 통해 이 특정 요청을 추적하기 위해 제공할 수 있는 문자열. 더 쉬운 디버깅과 모니터링을 위해 UUID 같은 임의 문자열을 사용하는 것이 매우 권장돼요. URL에 ?requestId=<your_uuid>처럼 추가해요.
  • beginMark(선택): 이전 insertReport 응답의 nextBeginMark 필드에 반환된 마커 값. 이 마커를 포함하면 반환되는 중복 이벤트 수를 줄여 후속 호출을 최적화하는 데 도움을 줘요. 참고: beginMark는 중복을 피하기 위한 힌트로 의도되었지만, 가끔 이벤트 반복이 여전히 발생할 수 있어요. beginMark가 지정되지 않으면 보고서는 지난 10분 동안의 수집 기록을 보여줘요. URL에 ?beginMark=<previous_nextBeginMark>처럼 추가해요.

요청 헤더:

  • Accept: 원하는 응답 형식을 지정. 허용 값은 text/plain 또는 application/json.
  • Authorization: Snowflake 인증 토큰. BEARER <jwt_token> 형식 사용.

요청 본문:

이 엔드포인트는 GET 요청에 대한 요청 본문을 받지 않아요. 필요한 매개 변수는 URL과 헤더에 제공돼요.

응답 본문:

응답 코드:

  • 200 — 성공. 보고서 반환.
  • 400 — 실패. 잘못된 형식 또는 한도 초과로 인한 유효하지 않은 요청.
  • 404 — 실패. pipeName이 인식되지 않음. 이 오류 코드는 엔드포인트를 호출할 때 사용된 역할에 충분한 권한이 없는 경우에도 반환될 수 있어요. 자세한 내용은 액세스 권한 부여를 참고해요.
  • 429 — 실패. 요청 속도 한도를 초과함.
  • 500 — 실패. 내부 오류 발생.

응답 페이로드:

성공 응답(200)은 최근에 테이블에 추가된 파일에 대한 정보를 포함해요. 이 보고서는 대용량 파일의 일부만 나타낼 수 있다는 점에 주의해요.

예를 들어:

{
  "pipe": "TESTDB.TESTSCHEMA.pipe2",
  "completeResult": true,
  "nextBeginMark": "1_39",
  "files": [
    {
      "path": "data2859002086815673867.csv",
      "stageLocation": "s3://mybucket/",
      "fileSize": 57,
      "timeReceived": "2017-06-21T04:47:41.453Z",
      "lastInsertTime": "2017-06-21T04:48:28.575Z",
      "rowsInserted": 1,
      "rowsParsed": 1,
      "errorsSeen": 0,
      "errorLimit": 1,
      "complete": true,
      "status": "LOADED"
    }
  ]
}

응답 필드:

필드 유형 설명
pipe String 파이프의 정규화된 이름.
completeResult Boolean 제공된 beginMark와 이 보고서 기록의 첫 번째 이벤트 사이에 이벤트가 놓치면 false. 그렇지 않으면 true.
nextBeginMark String 중복 레코드를 피하기 위해 다음 요청에 사용할 beginMark. 이 값은 힌트라는 점에 주의해요. 중복이 가끔 여전히 발생할 수 있어요.
files Array 기록 응답의 일부인 각 파일에 대한 하나의 JSON 객체 배열.
path String 스테이지 위치에 상대적인 파일 경로.
stageLocation String 파이프에 정의된 스테이지 ID(내부 스테이지) 또는 S3 버킷(외부 스테이지).
fileSize Long 파일 크기(바이트 단위).
timeReceived String 이 파일이 처리를 위해 수신된 시간. 형식은 UTC 시간대의 ISO-8601.
lastInsertTime String 이 파일의 데이터가 테이블에 마지막으로 삽입된 시간. 형식은 UTC 시간대의 ISO-8601.
rowsInserted Long 파일에서 대상 테이블에 삽입된 행 수.
rowsParsed Long 파일에서 파싱된 행 수. 오류가 있는 행은 건너뛸 수 있음.
errorsSeen Integer 파일에서 발생한 오류 수.
errorLimit Integer 파일이 실패한 것으로 간주되기 전에 허용되는 오류 수(ON_ERROR 복사 옵션 기준).
firstError [1] String 이 파일에서 만난 첫 번째 오류의 오류 메시지.
firstErrorLineNum [1] Long 첫 번째 오류의 줄 번호.
firstErrorCharacterPos [1] Long 첫 번째 오류의 문자 위치.
firstErrorColumnName [1] String 첫 번째 오류가 발생한 컬럼 이름.
systemError [1] String 파일이 처리되지 않은 이유를 설명하는 일반 오류.
complete Boolean 파일이 완전히 성공적으로 처리되었는지 여부.
status String 파일의 로드 상태: LOAD_IN_PROGRESS(파일 일부가 테이블에 로드되었지만 로드 프로세스가 아직 완료되지 않음), LOADED(전체 파일이 테이블에 로드됨), LOAD_FAILED(파일 로드 실패), PARTIALLY_LOADED(이 파일의 일부 행은 성공적으로 로드되었지만 다른 행은 오류로 로드되지 않음. 이 파일의 처리는 완료됨).

[1] 이 필드의 값은 파일에 오류가 포함된 경우에만 제공돼요.

엔드포인트: loadHistoryScan

내용이 테이블에 추가된 수집된 파일에 대한 보고서를 가져와요. 대용량 파일의 경우 이 보고서가 파일의 일부만일 수 있다는 점에 주의해요. 이 엔드포인트는 insertReport와 달리 두 시점 사이의 기록을 보는 점에서 차이가 있어요. 반환되는 항목은 최대 10,000개이지만, 원하는 기간을 커버하기 위해 여러 호출을 실행할 수 있어요.

중요 — 이 엔드포인트는 과도한 호출을 피하기 위해 속도가 제한돼요. 속도 한도(오류 코드 429)를 초과하지 않도록 loadHistoryScan보다 insertReport에 더 많이 의존할 것을 권장해요. loadHistoryScan을 호출할 때는 데이터 로드 집합을 포함하는 가장 좁은 시간 범위를 지정해요. 예를 들어 8분마다 지난 10분의 기록을 읽는 것은 잘 작동해요. 매분 지난 24시간의 기록을 읽으려고 하면 속도 한도에 도달했음을 나타내는 429 오류가 발생해요. 속도 한도는 각 기록 레코드를 몇 번 읽을 수 있도록 설계되었어요.

이러한 한도 없이 더 포괄적인 보기를 위해 Snowflake는 파이프 또는 테이블의 로드 기록을 반환하는 정보 스키마 테이블 함수 COPY_HISTORY를 제공해요.

메서드: GET

GET URL:

https://{account}.snowflakecomputing.com/v1/data/pipes/{pipeName}/loadHistoryScan?startTimeInclusive=<startTime>&endTimeExclusive=<endTime>&requestId=<requestId>

URL 매개 변수:

  • account(필수): Snowflake의 고유 계정 식별자.
  • pipeName(필수): 대소문자를 구분하는 정규화된 Snowpipe 이름. 예: myDatabase.mySchema.myPipe.
  • startTimeInclusive(필수): 로드 기록 데이터를 검색할 시간 범위의 시작. ISO-8601 형식의 타임스탬프로 지정(예: 2023-10-26T10:00:00Z). 이 타임스탬프는 쿼리의 포함 하한을 표시해요.
  • endTimeExclusive(선택): 로드 기록 데이터를 검색할 시간 범위의 끝. ISO-8601 형식의 타임스탬프로 지정(예: 2023-10-26T10:15:00Z). 이 타임스탬프는 쿼리의 배타적 상한을 표시해요. 이 매개 변수를 생략하면 현재 서버 타임스탬프(CURRENT_TIMESTAMP())가 시간 범위의 끝으로 사용돼요.
  • requestId(선택): Snowflake 시스템을 통해 이 특정 요청을 추적하기 위해 제공할 수 있는 문자열. 더 쉬운 디버깅과 모니터링을 위해 UUID 같은 임의 문자열을 사용할 것을 권장해요. URL에 ?requestId=<your_uuid>처럼 추가해요.

요청 헤더:

  • Accept: 원하는 응답 형식을 지정. 허용 값은 text/plain 또는 application/json.
  • Authorization: Snowflake 인증 토큰. BEARER <jwt_token> 형식 사용.

요청 본문:

이 엔드포인트는 GET 요청에 대한 요청 본문을 받지 않아요. 모든 필수 매개 변수는 URL과 헤더에 제공돼요.

응답 본문:

응답 코드:

  • 200 — 성공. 로드 기록 스캔 결과가 반환됨.
  • 400 — 실패. 잘못된 형식 또는 한도 초과로 인한 유효하지 않은 요청.
  • 404 — 실패. pipeName이 인식되지 않음.
  • 429 — 실패. 요청 속도 한도를 초과함.
  • 500 — 실패. 내부 오류 발생.

응답 페이로드:

성공 응답(200)은 최근에 테이블에 추가된 파일에 대한 정보를 포함해요. 이 보고서는 대용량 파일의 일부만 나타낼 수 있다는 점에 주의해요.

예를 들어:

{
  "pipe": "TESTDB.TESTSCHEMA.pipe2",
  "completeResult": true,
  "startTimeInclusive": "2017-08-25T18:42:31.081Z",
  "endTimeExclusive":"2017-08-25T22:43:45.552Z",
  "rangeStartTime":"2017-08-25T22:43:45.383Z",
  "rangeEndTime":"2017-08-25T22:43:45.383Z",
  "files": [
    {
      "path": "data2859002086815673867.csv",
      "stageLocation": "s3://mystage/",
      "fileSize": 57,
      "timeReceived": "2017-08-25T22:43:45.383Z",
      "lastInsertTime": "2017-08-25T22:43:45.383Z",
      "rowsInserted": 1,
      "rowsParsed": 1,
      "errorsSeen": 0,
      "errorLimit": 1,
      "complete": true,
      "status": "LOADED"
    }
  ]
}

응답 필드:

필드 유형 설명
pipe String 파이프의 정규화된 이름.
completeResult Boolean 보고서가 불완전하면(즉, 지정된 시간 범위의 항목 수가 10,000 항목 한도를 초과하면) false. false이면 사용자는 현재 rangeEndTime 값을 다음 요청의 startTimeInclusive 값으로 지정해 다음 항목 집합으로 진행할 수 있어요.
startTimeInclusive String 요청에 제공된 시작 타임스탬프(ISO-8601 형식).
endTimeExclusive String 요청에 제공된 종료 타임스탬프(ISO-8601 형식).
rangeStartTime String 응답에 포함된 파일에서 가장 오래된 항목의 타임스탬프(ISO-8601 형식).
rangeEndTime String 응답에 포함된 파일에서 가장 최근 항목의 타임스탬프(ISO-8601 형식).
files Array 기록 응답의 일부인 각 파일에 대한 하나의 JSON 객체 배열. 배열 안에서 응답 필드는 insertReport 응답에서 반환된 것과 동일해요.

더 알아보기