データ接続と統合Python仮想テーブルと計算プッシュダウンSnowflake の計算プッシュダウン

注: 以下の翻訳の正確性は検証されていません。AIPを利用して英語版の原文から機械的に翻訳されたものです。

Snowflake コンピュートプッシュダウン

Snowflake でコンピュートプッシュダウンを使用するには、Python リポジトリを作成し、transforms-tables ライブラリの最新バージョンをインストールします。

Snowpark ↗ セッションは、トランスフォームの入力や出力として設定された Snowflake テーブルの接続情報に基づいて設定されます。データは Snowpark DataFrame API を使用して加工できます。Snowpark API の詳しいガイダンスについては、Snowpark ドキュメント ↗を参照してください。

次の Snowpark トランスフォームの例では、Foundry Spark 構文ではなく Foundry 軽量 API 構文を使用していることに注意してください。

Copied!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 from snowflake.snowpark.functions import col, udf from snowflake.snowpark.types import StringType import snowflake.snowpark as snow from transforms.api import transform from transforms.tables import ( SnowflakeTable, TableInput, TableOutput, SnowflakeInput, SnowflakeOutput, ) ID_PREFIX = "CUSTOMER-NO-" @transform.snowflake.using( input_table=TableInput("ri.tables.main.table.1234"), output_table=TableOutput( "ri.tables.main.table.5678", "ri.magritte..source.1234", SnowflakeTable("DATABASE", "PUBLIC", "CUSTOMERS_CLEANED"), ), ) def compute_in_snowflake(input_table: SnowflakeInput, output_table: SnowflakeOutput): """ Snowflake テーブルでは、Snowpark API を使って軽量なトランスフォームを実行できます。これらの処理はすべて 基盤となる Snowflake インスタンスにプッシュダウンされるため、ビッグデータのワークロードに対応できます。このような構成では、 すべてのデータが同じ Snowpark インスタンス内に存在し、同じ接続を通じてアクセスできる必要があります。 """ # Snowpark DataFrame インスタンスを取得します df: snow.DataFrame = input_table.dataframe() session: snow.Session = df.session # データに適用する UDF を定義します @udf(session=session, return_type=StringType()) def fix_id_col(ident: int) -> str: """ ID を文字列に変換し、先頭に "CUSTOMER-NO-" を付加する UDF です。 """ return ID_PREFIX + str(ident) # UDF を適用します df = df.with_column("ID", fix_id_col(col("ID"))) # 新しいテーブルに書き戻します output_table.write(df)

コンピュート設定

デフォルトでは、Foundry はソースで設定したウェアハウスを使用してコンピュートリソースを割り当てます。

また、.with_warehouse(warehouse="MY_WAREHOUSE") を使用して、特定のウェアハウスへの接続を設定することもできます。

Copied!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 from transforms.api import transform from transforms.tables import ( SnowflakeTable, TableInput, TableOutput, SnowflakeInput, SnowflakeOutput, ) # カスタムウェアハウス設定を使用するトランスフォーム @transform.snowflake( source_table=TableInput( "ri.tables.main.table.my_snowflake_input_table" ), output_table=TableOutput( "ri.tables.main.table.my_snowflake_output_table", "ri.magritte..source.my_snowflake_source", SnowflakeTable("MY_DATABASE", "PUBLIC", "MY_TABLE"), ), ).with_warehouse("my_warehouse") def compute(source_table: SnowflakeInput, output_table: SnowflakeOutput): df: snow.DataFrame = source_table.dataframe() output_table.write(df)

Spark Connect で PySpark API を使用する

Snowpark の代わりに、.with_engine(engine="pyspark") を使用して Snowpark Connect セッションを確立できます。これにより、Spark Connect を通じて PySpark API を使用できるようになります。詳細については、Snowpark Connect ↗を参照してください。

Copied!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 from pyspark.sql.connect.dataframe import DataFrame from pyspark.sql.connect.functions import lit from pyspark.sql.connect.session import SparkSession from transforms.api import transform from transforms.tables import ( TableInput, TableOutput, SnowparkConnectInput, SnowparkConnectOutput, SnowflakeTable ) @transform.snowflake.using( input_table=TableInput("ri.tables.main.table.1234"), output_table=TableOutput( "ri.tables.main.table.5678", "ri.magritte..source.1234", SnowflakeTable("DATABASE", "PUBLIC", "CUSTOMERS_CLEANED_ANON"), ), ).with_engine("pyspark") def compute(source_table: SnowparkConnectInput, output_table: SnowparkConnectOutput): # これは Snowpark ではなく Spark Connect の DataFrame になります df: DataFrame = source_table.dataframe() # 同様に、これは Spark Connect のセッションになります session: SparkSession = source_table.session df = df.withColumn("test", lit(3)) output_table.write_dataframe(df)

データを pandas DataFrame に変換する

Snowpark API では、データを pandas DataFrame に変換できます。データの規模が十分に小さい場合は、この方法を使用して Snowflake から Foundry 軽量コンピュートにデータを取り込めます。これにより、Snowpark API の機能を超えるトランスフォームを使用でき、Snowflake テーブルをほかの Foundry データと組み合わせることができます。

Copied!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 import hashlib from transforms.api import transform from transforms.tables import ( SnowflakeTable, TableInput, TableOutput, SnowflakeInput, SnowflakeOutput, ) @transform.snowflake.using( input_table=TableInput("ri.tables.main.table.1234"), output_table=TableOutput( "ri.tables.main.table.5678", "ri.magritte..source.1234", SnowflakeTable("DATABASE", "PUBLIC", "CUSTOMERS_CLEANED_ANON"), ), ) def compute_local(input_table: SnowflakeInput, output_table: SnowflakeOutput): """ Snowpark は pandas DataFrame への変換もサポートしているため、Snowflake テーブルに対して軽量な トランスフォームを使用し、コンテナ内で計算処理を実行できます。この機能を使うと、 Snowpark でサポートされている範囲を超えた処理を実行できます。 """ # Snowpark DataFrame インスタンスを取得します df = input_table.dataframe() session = df.session # pandas に変換します pd_df = df.to_pandas() # CITY、STATE、ZIP_CODE を連結した値をハッシュ化して ANON_CODE を作成します def generate_anon_code(row): concatenated = f"{row['CITY']}{row['STATE']}{row['ZIP_CODE']}" return hashlib.sha256(concatenated.encode("utf-8")).hexdigest() # 関数を適用して ANON_CODE 列を作成します pd_df["ANON_CODE"] = pd_df.apply(generate_anon_code, axis=1) # ID 列と ANON_CODE 列を選択します result_data = pd_df[["ID", "ANON_CODE"]] # 新しいテーブルに書き戻します new_df = session.create_dataframe(result_data, schema=["ID", "ANON_CODE"]) output_table.write(new_df)

Snowflake へのコンピュートプッシュダウンを使用する場合、@incremental デコレーターを使用したインクリメンタル(差分処理)計算は現在サポートされていません。