디렉터리 테이블을 사용한 데이터 처리 파이프라인 구축

디렉터리 테이블을 사용한 데이터 처리 파이프라인 구축

스테이지에서 파일 수준 메타데이터를 추적하고 저장하는 디렉터리 테이블을 스트림과 태스크 같은 다른 Snowflake 객체와 결합해 데이터 처리 파이프라인을 구축해요.

스트림은 디렉터리 테이블, 테이블, 외부 테이블 또는 뷰의 기본 테이블에 대해 이루어진 DML(데이터 조작 언어) 변경을 기록해요. 태스크는 SQL 명령 또는 광범위한 사용자 정의 함수(UDF)가 될 수 있는 단일 작업을 실행해요. 태스크를 주기적으로 실행되도록 예약하거나, 온디맨드로 실행할 수 있어요.

출처: Documentation

본문

예제: PDF를 처리하는 간단한 파이프라인 만들기

이 예제는 다음을 수행하는 간단한 데이터 처리 파이프라인을 구축해요.

  • 스테이지에 추가된 PDF 파일 감지.
  • 파일에서 데이터 추출.
  • 데이터를 Snowflake 테이블에 삽입.

파이프라인은 스트림을 사용해 스테이지의 디렉터리 테이블에 대한 변경을 감지하고, UDF를 실행해 파일에서 데이터를 추출하는 태스크를 사용해요.

1단계: 디렉터리 테이블이 활성화된 스테이지 만들기

디렉터리 테이블이 활성화된 내부 스테이지를 만들어요. 예제 문은 ENCRYPTION 유형을 SNOWFLAKE_SSE로 설정해 스테이지에서 비정형 데이터 접근을 활성화해요.

CREATE OR REPLACE STAGE my_pdf_stage
  ENCRYPTION = ( TYPE = 'SNOWFLAKE_SSE')
  DIRECTORY = ( ENABLE = TRUE);

2단계: 디렉터리 테이블에 스트림 만들기

디렉터리 테이블이 속한 스테이지를 지정해 디렉터리 테이블에 스트림을 만들어요. 스트림은 디렉터리 테이블의 변경을 추적해요. 이 예제의 5단계에서 이 스트림을 사용해 태스크를 구성해요.

CREATE STREAM my_pdf_stream ON STAGE my_pdf_stage;

3단계: PDF를 파싱하는 사용자 정의 함수 만들기

PDF 파일에서 데이터를 추출하는 UDF를 만들어요. 나중 단계에서 만드는 태스크가 이 UDF를 호출해 스테이지에 새로 추가된 파일을 처리해요.

다음 예제 문은 제품 리뷰 데이터가 포함된 PDF 파일을 처리하는 PDF_PARSE라는 Python UDF를 만들어요. UDF는 PyPDF2 라이브러리를 사용해 양식 필드 데이터를 추출해요. 양식 이름과 값을 키-값 쌍으로 포함하는 딕셔너리를 반환해요.

참고 — UDF는 SnowflakeFile 클래스를 사용해 동적으로 지정된 파일을 읽어요. SnowflakeFile에 대한 자세한 내용은 SnowflakeFile로 동적으로 지정된 파일 읽기를 참고해요.

CREATE OR REPLACE FUNCTION PDF_PARSE(file_path string)
  RETURNS VARIANT
  LANGUAGE PYTHON
  RUNTIME_VERSION = '3.8'
  HANDLER = 'parse_pdf_fields'
  PACKAGES=('typing-extensions','PyPDF2','snowflake-snowpark-python')
AS
$$
from pathlib import Path
import PyPDF2 as pypdf
from io import BytesIO
from snowflake.snowpark.files import SnowflakeFile

def parse_pdf_fields(file_path):
    with SnowflakeFile.open(file_path, 'rb') as f:
        buffer = BytesIO(f.readall())
    reader = pypdf.PdfFileReader(buffer)
    fields = reader.getFields()
    field_dict = {}
    for k, v in fields.items():
        if "/V" in v.keys():
            field_dict[v["/T"]] = v["/V"].replace("/", "") if v["/V"].startswith("/") else v["/V"]

    return field_dict
$$;

4단계: 파일 내용을 저장할 테이블 만들기

다음으로 각 행이 file_name과 file_data라는 컬럼에 스테이지의 파일에 대한 정보를 저장하는 테이블을 만들어요. 나중 단계에서 만드는 태스크가 이 테이블에 데이터를 로드해요.

CREATE OR REPLACE TABLE prod_reviews (
  file_name varchar,
  file_data variant
);

5단계: 태스크 만들기

스트림에서 스테이지의 새 파일을 확인하고 파일 데이터를 prod_reviews 테이블에 삽입하는 예약된 태스크를 만들어요.

다음 문은 이전에 만든 스트림을 사용해 예약된 태스크를 만들어요. 태스크는 SYSTEM$STREAM_HAS_DATA 함수를 사용해 스트림에 CDC(변경 데이터 캡처) 레코드가 있는지 확인해요.

CREATE OR REPLACE TASK load_new_file_data
  WAREHOUSE = 'MY_WAREHOUSE'
  SCHEDULE = '1 minute'
  COMMENT = 'Process new files on the stage and insert their data into the prod_reviews table.'
  WHEN
  SYSTEM$STREAM_HAS_DATA('my_pdf_stream')
  AS
  INSERT INTO prod_reviews (
    SELECT relative_path as file_name,
    PDF_PARSE(build_scoped_file_url('@my_pdf_stage', relative_path)) as file_data
    FROM my_pdf_stream
    WHERE METADATA$ACTION='INSERT'
  );

6단계: 파이프라인을 테스트하기 위해 태스크 실행

파이프라인이 작동하는지 확인하려면 스테이지에 파일을 추가하고, 태스크를 수동으로 실행한 다음, product_reviews 테이블을 조회하면 돼요.

먼저 my_pdf_stage 스테이지에 PDF 파일 몇 개를 추가하고 스테이지를 새로 고쳐요.

참고 — 이 예제는 Snowflake 웹 인터페이스의 워크시트에서 실행할 수 없는 PUT 명령을 사용해요. Snowsight로 파일을 업로드하려면 이름이 있는 내부 스테이지에 파일 업로드를 참고해요.

PUT file:///my/file/path/prod_review1.pdf @my_pdf_stage AUTO_COMPRESS = FALSE;
PUT file:///my/file/path/prod_review2.pdf @my_pdf_stage AUTO_COMPRESS = FALSE;

ALTER STAGE my_pdf_stage REFRESH;

스트림을 조회해 스테이지에 추가한 두 PDF 파일을 기록했는지 확인할 수 있어요.

SELECT * FROM my_pdf_stream;

이제 태스크를 실행해 PDF 파일을 처리하고 product_reviews 테이블을 업데이트해요.

EXECUTE TASK load_new_file_data;
+----------------------------------------------------------+
| status                                                   |
|----------------------------------------------------------|
| Task LOAD_NEW_FILE_DATA is scheduled to run immediately. |
+----------------------------------------------------------+
1 Row(s) produced. Time Elapsed: 0.178s

product_reviews 테이블을 조회해 태스크가 각 PDF 파일에 대해 행을 추가했는지 확인해요.

select * from prod_reviews;
+------------------+----------------------------------+
| FILE_NAME        | FILE_DATA                        |
|------------------+----------------------------------|
| prod_review1.pdf | {                                |
|                  |   "FirstName": "John",           |
|                  |   "LastName": "Johnson",         |
|                  |   "Middle Name": "Michael",      |
|                  |   "Product": "Tennis Shoes",     |
|                  |   "Purchase Date": "03/15/2022", |
|                  |   "Recommend": "Yes"             |
|                  | }                                |
| prod_review2.pdf | {                                |
|                  |   "FirstName": "Emily",          |
|                  |   "LastName": "Smith",           |
|                  |   "Middle Name": "Ann",          |
|                  |   "Product": "Red Skateboard",   |
|                  |   "Purchase Date": "01/10/2023", |
|                  |   "Recommend": "MayBe"           |
|                  | }                                |
+------------------+----------------------------------+

마지막으로 FILE_DATA 컬럼의 객체를 개별 컬럼으로 파싱하는 뷰를 만들 수 있어요. 그런 다음 뷰를 조회해 파일 내용을 분석하고 작업할 수 있어요.

CREATE OR REPLACE VIEW prod_review_info_v
  AS
  WITH file_data
  AS (
      SELECT
        file_name
        , parse_json(file_data) AS file_data
      FROM prod_reviews
  )
  SELECT
      file_name
      , file_data:FirstName::varchar AS first_name
      , file_data:LastName::varchar AS last_name
      , file_data:"Middle Name"::varchar AS middle_name
      , file_data:Product::varchar AS product
      , file_data:"Purchase Date"::date AS purchase_date
      , file_data:Recommend::varchar AS recommended
      , build_scoped_file_url(@my_pdf_stage, file_name) AS scoped_review_url
  FROM file_data;

SELECT * FROM prod_review_info_v;

+------------------+------------+-----------+-------------+----------------+---------------+-------------+--------------...+
| FILE_NAME        | FIRST_NAME | LAST_NAME | MIDDLE_NAME | PRODUCT        | PURCHASE_DATE | RECOMMENDED | SCOPED_REVIEW_URL      |
|------------------+------------+-----------+-------------+----------------+---------------+-------------+----------------------+
| prod_review1.pdf | John       | Johnson   | Michael     | Tennis Shoes   | 2022-03-15    | Yes         | https://...           |
| prod_review2.pdf | Emily      | Smith     | Ann         | Red Skateboard | 2023-01-10    | MayBe       | https://...           |
+------------------+------------+-----------+-------------+----------------+---------------+-------------+----------------------+

더 알아보기