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 응답에서 반환된 것과 동일해요. |