For AI agents: the complete documentation index is available at https://docs.dataplatform.ovh.net/ja/llms.txt, the full documentation bundle is available at https://docs.dataplatform.ovh.net/ja/llms-full.txt, and this page is available as Markdown at https://docs.dataplatform.ovh.net/ja/tutorials-sql-transformation.md.
  • 🇯🇵 日本語
  • OVHcloud Data Platform内でのSQL変換

    SQL(構造化問い合わせ言語)は、構造化データの問い合わせと変換に最も広く使用され、効率的な言語です

    目的

    OVHcloud Data Platform内でのSQL変換のガイドへようこそ。SQL(構造化問い合わせ言語)は、構造化データの問い合わせと変換に最も広く使用され、効率的な言語です。このドキュメントでは、データ操作と集約のためのSQLの基本概念を紹介し、プラットフォーム内でその強力な機能を活用する方法を説明します。

    大規模なデータセットをクリーンアップ、再構成、フィルタリング、または集約する必要がある場合、SQLはデータ変換の目標を達成するための堅牢で直感的な方法を提供します。その宣言的な性質により、データ操作の方法ではなく、データで達成したいことを中心に考えることができ、さまざまなデータ専門家にとってアクセスしやすくなります。

    データ変換のためのSQLの理由

    SQLは、データ変換パイプラインにおいていくつかの重要な理由から不可欠なツールです:

    • 普遍性: 標準言語であり、さまざまなデータベースやデータプラットフォームで広く採用されており、スキルを簡単に移行できます。
    • 可読性と簡潔さ: 英語のような構文により、複雑な操作を比較的簡単に理解して記述できます。
    • パフォーマンス: SQLエンジンは関係操作に高度に最適化されており、大規模なデータ操作においてカスタムコードを上回ることがよくあります。
    • 宣言的な性質: データの最終状態を指定し、エンジンが最も効率的な方法でそれを達成する方法を決定します。
    • 統合: データウェアハウス、ビジネスインテリジェンス、レポートツールとシームレスに統合されます。

    SQL変換の基本概念

    SQL変換の核心は、標準的なSQLコマンドを使用してデータを操作することです。これは次のようなものを含むことがあります。

    1. データのクリーンアップと準備

    • フィルタリング: WHERE 句を使用して、条件に基づいて特定の行を選択します。
    • 選択/投影: SELECT 文を使用して特定の列を選択し、名前を変更します(AS)。
    • 型変換: CAST または TRY_CAST 関数を使用してデータ型を変更します。TRY_CAST は、Trinoで変換エラーを優雅に処理する(NULL を返す代わりにクラッシュする)ために特に便利です。
    • 欠損値の処理: COALESCE または CASE 文を使用して NULL 値を置き換えます。
    • 文字列操作: SUBSTRINGLENGTHUPPERLOWERTRIM のような関数を使用してテキストデータをクリーンアップし、標準化します。

    2. データの集約と要約

    • グループ化: GROUP BY を使用して、指定された列で同じ値を持つ行を集約します。
    • 集約関数: COUNTSUMAVGMINMAX のような関数を適用して集約されたグループを要約します。
    • 集約のフィルタリング: HAVING 句を使用して GROUP BY 操作の結果をフィルタリングします。

    3. データの再構成と再構造化

    • 結合: 関連する列に基づいて2つ以上のテーブルのデータを結合します(INNER JOINLEFT JOINRIGHT JOINFULL OUTER JOIN)。
    • ユニオン: 2つ以上の SELECT 文の結果セットを結合します(UNIONUNION ALL)。
    • ピボット/アンピボット: 行を列(ピボット)または列を行(アンピボット)に変換してデータの構造を変更します。Trino SQL(DPEで一般的に使用される)では、ピボットは SUMFILTER または CASE 文を使用してよく実現され、アンピボットは CROSS JOIN UNNEST を使用して実現されます。
    • ウィンドウ関数: 現在の行に関連するテーブル行のセット全体で計算を実行し、行を折りたたむことなく(ROW_NUMBER()RANK()LEAD()LAG()SUM() OVER()AVG() OVER())。

    OVHcloud Data Platform 上での SQL 変換(実践ガイド)

    このセクションでは、OVHcloud Data Platform の DPE(Data Processing Environment)ノートブック内で SQL 変換を実行する、実践的なエンドツーエンドの例を提供します。具体的には、SDK を使用し、dirty_cafe_sales データセットに焦点を当てます。

    前提条件

    1. dirty_cafe_sales.csv をダウンロードします。
    2. Connectors にアップロードし、Analyzer を使用してメタデータを抽出します。
    3. Lakehouse ManagerTable をソースから新規作成します。

    データセット情報:

    私たちは、dirty_cafe_sales という名前のシミュレートされたカフェ販売データセットを使用します。このデータセットには、次の列と既知のデータ品質の問題が含まれています。

    列名説明既知の問題
    transaction_id各取引のユニーク識別子なし
    item販売されたアイテムの名前なし
    quantity取引で販売された単位の数'UNKNOWN' 文字列を含む
    price_per_unitアイテムの単価なし
    total_spent取引のアイテム行の合計支払額'ERROR' 文字列を含む
    payment_method支払い方法(例:クレジットカード、現金)なし
    location販売場所(例:店内、テイクアウト)'ERROR' 文字列と空文字列を含む ''
    transaction_date取引日非日付文字列(例:'ERROR')を含む。キャストとエラー処理が必要

    1. 新しい DPE ノートブックを作成する

    • Data Processing Environment (DPE) → Notebooks に移動します。
    • + New Notebook をクリックし、このガイドの目的で Base Notebook を続行します。
    • ノートブックに意味のある名前を付けます(例:Cafe_Sales_SQL_Transformations)。
    • Create をクリックします。
    • JupyterLab が開いたら、Python3 ノートブックをクリックして新しい .ipynb ファイルを作成します。ここが、次のすべての手順をセルで実行する場所です。

    2. SDK に接続する

    DPE 環境で提供される SDK を使用します。この SDK は、Lakehouse Manager に接続し、テーブルとやり取りする方法を提供します。

    # Import necessary modules from SDK
    from forepaas.dwh import connect, bulk_insert
    from forepaas.core.settings import CONFIG
    from forepaas.dwh.logical import LogicalObject
    import pandas as pd #useful for displaying dataframes
    
    print("SDK modules imported successfully.")

    3. データセットからテーブルをリストする

    dirty_cafe_sales に接続する前に、指定されたデータパス内のテーブルをリストすることは、その存在と正確な名前を確認するための良い習慣です。

    # Connect to the Lakehouse Manager
    # Connect to the default Lakehouse Manager dataset
    connector = connect("dwh/default_dataset/")
    
    print("Listing tables in 'dwh/default_dataset/':")
    available_tables = connector.list()
    for table_name in available_tables:
        print(f"- {table_name}")
    
    if "dirty_cafe_sales" in available_tables:
        print("\n'dirty_cafe_sales' table found!")
    else:
        print("\nWARNING: 'dirty_cafe_sales' table not found. Please check the table name or path.")

    4. テーブルに接続してデータを検査する

    次に、dirty_cafe_sales テーブルに connector.select() を使用して接続し、その情報と記述統計を印刷します。このステップでは、データを視覚的に確認できます。これは、クリーンアップが必要な 'ERROR' と 'UNKNOWN' 値を含みます。

    # Connect to the 'dirty_cafe_sales' table
    df_raw_sales = connector.select("dirty_cafe_sales")
    
    print("\nDataset information:")
    df_raw_sales.info()
    
    print("\nDataset details:")
    df_raw_sales.describe()
    
    print("\nSample of Raw Data (first 5 rows):")
    display(df_raw_sales.head())

    5. SQL コマンド(変換)を実行する

    これが変換の核心です。基本的な探索用のシンプルな SQL クエリと、徹底的なクリーンアップと詳細な集計用のより複雑な SQL クエリの 2 つを定義します。

    Info

    セミコロンに関する重要な注意点: ノートブックのようなプログラム環境で SDK または API を介して SQL を実行する場合、SQL クエリ文字列の末尾にトレーリングセミコロン (;) を含めないでください。API は通常、末尾に明示的な区切り文字のない単一の SQL ステートメントを期待します。これを含めると、mismatched input ';'syntax error near ';' のような一般的なエラーが発生することがあります。

    5.1 シンプルな SQL クエリ: 場所別の日別総販売額(初期探索)

    このクエリは、基本的な集計を示し、最初に生データの問題がエラーを引き起こす方法を示し、その後 TRY_CAST を使用して強固に修正します。また、location フィールドをクリーンアップします。

    print("--- Running Simple SQL Query ---")
    
    SIMPLE_SQL_QUERY = """
    SELECT
        CAST(valid_transaction_date AS DATE) AS sale_date,
        -- Handle 'ERROR' and empty strings in location
        CASE
            WHEN location = 'ERROR' THEN 'Unknown'
            WHEN TRIM(location) = '' THEN 'Unknown'
            ELSE location
        END AS clean_location,
        SUM(TRY_CAST(total_spent AS DOUBLE)) AS gross_revenue_dirty
    FROM
        (
            SELECT
                TRY_CAST(transaction_date AS DATE) AS valid_transaction_date,
                location,
                total_spent
            FROM
                dirty_cafe_sales
        ) AS subquery_sales
    WHERE
        valid_transaction_date IS NOT NULL
    GROUP BY
        CAST(valid_transaction_date AS DATE),
        CASE
            WHEN location = 'ERROR' THEN 'Unknown'
            WHEN TRIM(location) = '' THEN 'Unknown'
            ELSE location
        END
    ORDER BY
        sale_date DESC, clean_location
    """ # No semicolon at the end here!
    
    try:
        # Execute the query using your connector's method.
        df_simple_result = connector.query(SIMPLE_SQL_QUERY)
        
        # --- Pandas Post-Processing for Data Types ---
        # Convert 'sale_date' to datetime objects for proper date operations
        df_simple_result['sale_date'] = pd.to_datetime(df_simple_result['sale_date'])
        # --- End Pandas Post-Processing ---
    
        print("\nSimple Query Results (first 10 rows):")
        display(df_simple_result.head(10))
        print("\nSimple Query Results Schema:")
        df_simple_result.info()
    
    except Exception as e:
        print(f"Error executing Simple SQL query: {e}")
        print("\nFailed SQL Query:\n", SIMPLE_SQL_QUERY)

    5.2 複雑な SQL クエリ: 詳細なクリーンアップ済みの日別アイテムパフォーマンス

    このクエリは、データを強固にクリーンアップし、正確な販売メトリクスを計算し、日付、アイテム、場所別に集計します。これは、quantityUNKNOWNtotal_spentERROR、悪い transaction_date 文字列、および問題のある location エントリに特定されたデータ品質の問題に直接対応します。

    print("\n--- Running Complex SQL Transformation Query ---")
    
    COMPLEX_SQL_TRANSFORMATION_QUERY = """
    WITH cleaned_and_corrected_sales AS (
        SELECT
            transaction_id,
            item,
            -- Clean and cast quantity: 'UNKNOWN' becomes NULL, then cast to INTEGER.
            -- TRY_CAST handles non-numeric strings safely by returning NULL.
            TRY_CAST(NULLIF(quantity, 'UNKNOWN') AS INTEGER) AS quantity_cleaned,
            -- Ensure price_per_unit is numeric, handling potential non-numeric entries safely.
            TRY_CAST(price_per_unit AS DOUBLE) AS price_per_unit_cleaned,
            payment_method,
            -- Clean location: 'ERROR' and empty strings become 'Unknown'.
            CASE
                WHEN location = 'ERROR' THEN 'Unknown'
                WHEN TRIM(location) = '' THEN 'Unknown'
                ELSE location
            END AS clean_location,
            -- Use TRY_CAST for transaction_date to handle bad date strings safely, then filter later.
            TRY_CAST(transaction_date AS DATE) AS sale_date_raw
        FROM
            dirty_cafe_sales
    ),
    final_calculated_sales AS (
        SELECT
            transaction_id,
            item,
            quantity_cleaned AS final_quantity,
            price_per_unit_cleaned AS final_price_per_unit,
            -- Recalculate total_spent based on cleaned quantity and price_per_unit.
            quantity_cleaned * price_per_unit_cleaned AS calculated_total_spent,
            payment_method,
            clean_location,
            sale_date_raw AS sale_date
        FROM
            cleaned_and_corrected_sales
        -- Filter out rows where crucial values (quantity, price_per_unit, or sale_date)
        -- couldn't be cleanly converted, ensuring only valid data proceeds.
        WHERE
            quantity_cleaned IS NOT NULL 
            AND price_per_unit_cleaned IS NOT NULL
            AND sale_date_raw IS NOT NULL 
    )
    SELECT
        sale_date,
        item,
        clean_location AS location,
        COUNT(DISTINCT transaction_id) AS number_of_transactions,
        SUM(final_quantity) AS total_items_sold,
        SUM(calculated_total_spent) AS total_revenue_cleaned,
        AVG(final_price_per_unit) AS average_item_price_per_unit,
        -- Pivot revenue by payment method using Trino's FILTER clause
        SUM(calculated_total_spent) FILTER (WHERE payment_method = 'Credit Card') AS revenue_credit_card,
        SUM(calculated_total_spent) FILTER (WHERE payment_method = 'Cash') AS revenue_cash,
        SUM(calculated_total_spent) FILTER (WHERE payment_method = 'Digital Wallet') AS revenue_digital_wallet,
        -- Calculate average quantity per transaction for this group
        CAST(SUM(final_quantity) AS DOUBLE) / CAST(COUNT(DISTINCT transaction_id) AS DOUBLE) AS avg_quantity_per_transaction
    FROM
        final_calculated_sales
    GROUP BY
        sale_date,
        item,
        clean_location
    ORDER BY
        sale_date DESC, total_revenue_cleaned DESC
    """ # No semicolon at the end here!
    
    try:
        # Execute the complex transformation query
        df_transformed = connector.query(COMPLEX_SQL_TRANSFORMATION_QUERY)
    
        # --- Pandas Post-Processing for Data Types ---
        # Convert 'sale_date' to datetime objects for proper date operations
        df_transformed['sale_date'] = pd.to_datetime(df_transformed['sale_date'])
        # --- End Pandas Post-Processing ---
    
        print("\nTransformed Cafe Sales Data (Sample - first 5 rows):")
        display(df_transformed.head())
        print(f"\nTransformed Data Schema:")
        df_transformed.info()
    
        print(f"\nTotal rows in transformed data: {len(df_transformed)}")
    
    except Exception as e:
        print(f"Error during Complex SQL transformation: {e}")
        print("\nFailed SQL Query:\n", COMPLEX_SQL_TRANSFORMATION_QUERY)

    6. 変換されたデータから物理テーブルを作成する(CTAS)

    複雑な変換を成功裏に実行し、結果を確認した後、次の論理的なステップは、データベース内の新しい物理テーブルにこのクリーンアップ済みの集計データを保持することです。これは、通常 CREATE TABLE AS SELECT(CTAS)ステートメントを使用して実行されます。この新しいテーブルは、レポート作成、さらに分析、または他のデータ処理のソースとして使用できます。複雑なクリーンアップロジックを毎回再実行する必要はありません。

    Info

    SQL 実行に関する重要な注意点: 一部のデータベースコネクタまたは API は、query() 呼び出しごとに 1 つの SQL ステートメントのみを期待します。DROP TABLECREATE TABLE AS SELECT を実行するには、それらを別々のコマンドとして送信します。また、各クエリ文字列の末尾にトレーリングセミコロンがないことを確認してください。

    print("\n--- Creating Physical Table from Complex SQL Transformation Query ---")
    
    # Define the name of your new cleaned table. This variable can be reused across cells.
    NEW_CLEANED_TABLE_NAME = "cleaned_cafe_sales_daily_summary"
    
    # 1. DROP TABLE statement (removes the table if it already exists, for idempotent runs)
    # Note: No trailing semicolon at the very end of the string.
    DROP_TABLE_QUERY = f"DROP TABLE IF EXISTS {NEW_CLEANED_TABLE_NAME}"
    
    # 2. CREATE TABLE AS SELECT statement
    # This uses the same logic from the COMPLEX_SQL_TRANSFORMATION_QUERY
    # Note: No trailing semicolon at the very end of the string.
    CTAS_CORE_QUERY = f"""
    CREATE TABLE {NEW_CLEANED_TABLE_NAME} AS
    WITH cleaned_and_corrected_sales AS (
        SELECT
            transaction_id,
            item,
            TRY_CAST(NULLIF(quantity, 'UNKNOWN') AS INTEGER) AS quantity_cleaned,
            TRY_CAST(price_per_unit AS DOUBLE) AS price_per_unit_cleaned,
            payment_method,
            CASE
                WHEN location = 'ERROR' THEN 'Unknown'
                WHEN TRIM(location) = '' THEN 'Unknown'
                ELSE location
            END AS clean_location,
            TRY_CAST(transaction_date AS DATE) AS sale_date_raw
        FROM
            dirty_cafe_sales
    ),
    final_calculated_sales AS (
        SELECT
            transaction_id,
            item,
            quantity_cleaned AS final_quantity,
            price_per_unit_cleaned AS final_price_per_unit,
            quantity_cleaned * price_per_unit_cleaned AS calculated_total_spent,
            payment_method,
            clean_location,
            sale_date_raw AS sale_date
        FROM
            cleaned_and_corrected_sales
        WHERE
            quantity_cleaned IS NOT NULL 
            AND price_per_unit_cleaned IS NOT NULL
            AND sale_date_raw IS NOT NULL 
    )
    SELECT
        sale_date,
        item,
        clean_location AS location,
        COUNT(DISTINCT transaction_id) AS number_of_transactions,
        SUM(final_quantity) AS total_items_sold,
        SUM(calculated_total_spent) AS total_revenue_cleaned,
        AVG(final_price_per_unit) AS average_item_price_per_unit,
        SUM(calculated_total_spent) FILTER (WHERE payment_method = 'Credit Card') AS revenue_credit_card,
        SUM(calculated_total_spent) FILTER (WHERE payment_method = 'Cash') AS revenue_cash,
        SUM(calculated_total_spent) FILTER (WHERE payment_method = 'Digital Wallet') AS revenue_digital_wallet,
        CAST(SUM(final_quantity) AS DOUBLE) / CAST(COUNT(DISTINCT transaction_id) AS DOUBLE) AS avg_quantity_per_transaction
    FROM
        final_calculated_sales
    GROUP BY
        sale_date,
        item,
        clean_location
    ORDER BY
        sale_date DESC, total_revenue_cleaned DESC
    """
    
    try:
        # Execute DROP TABLE first
        print(f"Dropping table {NEW_CLEANED_TABLE_NAME} if it exists...")
        connector.query(DROP_TABLE_QUERY)
        print("Drop table command executed.")
    
        # Then execute CREATE TABLE AS SELECT
        print(f"Creating table {NEW_CLEANED_TABLE_NAME}...")
        # For DDL operations like CREATE TABLE, connector.query() might return an empty DataFrame or None.
        connector.query(CTAS_CORE_QUERY)
        
        print(f"\nSuccessfully created table: {NEW_CLEANED_TABLE_NAME}")
        
        # Optional: Verify the table was created by querying its schema or a few rows
        print(f"\nVerifying schema of new table: {NEW_CLEANED_TABLE_NAME}")
        df_verify = connector.query(f"SELECT * FROM {NEW_CLEANED_TABLE_NAME} LIMIT 5")
        display(df_verify)
        df_verify.info()
    
    except Exception as e:
        print(f"Error during CTAS operation for {NEW_CLEANED_TABLE_NAME}: {e}")
        if "DROP TABLE" in str(e) and DROP_TABLE_QUERY in str(e):
            print("\nFailed SQL Query (DROP TABLE):\n", DROP_TABLE_QUERY)
        elif "CREATE TABLE" in str(e) and CTAS_CORE_QUERY in str(e):
            print("\nFailed SQL Query (CREATE TABLE AS SELECT):\n", CTAS_CORE_QUERY)
        else:
            print("\nFailed SQL Query:\n", e)

    7. 新しい物理テーブルを探索する

    標準的な SQL メタデータコマンドを使用して、新しく作成した物理テーブルを探索できます。<your_catalog_name><your_schema_name> を実際の値に置き換えることを忘れないでください。

    # You'll need to know your catalog and schema names.
    # Example: your_catalog_name = "default_dataset", your_schema_name = "sales_data"
    # To retrieve catalog and schema name you can execute the following commands in the Lakehouse Manager Explorer
    # CATALOG LIST: show catalogs
    # SCHEMA LIST: show schemas from {catalog_name}
    
    your_catalog_name = "your_main_catalog" # <<< IMPORTANT: Replace with your actual catalog name!
    your_schema_name = "your_schema_name"   # <<< IMPORTANT: Replace with your actual schema name!
    your_table_name = "cleaned_cafe_sales_daily_summary"
    
    print(f"\n--- Exploring the New Table: {your_table_name} ---")
    
    try:
        # Describe your new table
        df_describe = connector.query(f"DESCRIBE {your_catalog_name}.{your_schema_name}.{your_table_name}")
        print(f"\nDescription of table '{your_table_name}':")
        display(df_describe)
        
        # Select some data from your new table
        df_sample_data = connector.query(f"SELECT * FROM {your_catalog_name}.{your_schema_name}.{your_table_name} LIMIT 5")
        print(f"\nSample data from '{your_table_name}':")
        display(df_sample_data)
    
    except Exception as e:
        print(f"Error exploring table '{your_table_name}': {e}")
        print("Please check the full table path (catalog.schema.table) and permissions.")

    8. 物理テーブルから論理オブジェクトを作成する

    SQL を介して新しいテーブルを作成すると、そのテーブルはデータベース内の物理テーブルとしてのみ存在します。ただし、UI の Tables セクションは論理レベルで動作します。これは、プラットフォーム内の論理オブジェクトとして登録されたテーブルを表示します。

    新しく作成した物理テーブルを UI で表示し、使用できるようにするには、対応する論理オブジェクトを作成する必要があります。この論理的な表現は、データベースとプラットフォームインターフェイスの間のブリッジとして機能し、UI から直接テーブルのスキーマとデータとやり取りできるようにします。

    print(f"\n--- Creating LogicalObject for '{NEW_CLEANED_TABLE_NAME}' ---")
    
    try:
            
        logical_cleaned_sales = LogicalObject().create_from_physical(NEW_CLEANED_TABLE_NAME)
        print(f"Successfully created logical object for: {NEW_CLEANED_TABLE_NAME}")
    
    except NameError:
        print("Error: 'LogicalObject' is not defined. Please ensure you have imported the correct library/class or that it's globally accessible.")
    except Exception as e:
        print(f"Error creating logical object: {e}")
    Info

    論理オブジェクトを削除する場合は、UI からテーブルを直接削除するか、LogicalObject().remove("table_name") メソッドを使用できます。これは、論理レベルと物理レベルの両方でテーブルを削除します。

    結論

    OVHcloud Data Platform 上で SQL 変換プロセスをエンドツーエンドで実行する方法を成功裏にナビゲートしました。生の汚れたデータセットから始まり、さまざまな SQL クリーンアップ、集計、および再形成技術を適用して、クリーンで要約された高度に使用可能なデータセットを生成しました。この変換されたデータは、新しい物理テーブルに保持され、Python ワークフローへのシームレスな統合のために論理オブジェクトとして表現されました。この基礎的な知識は、より複雑なデータ準備の課題に対処し、強固なデータパイプラインを構築するための力を与えてくれます。

    さらに進む

    ソリューションを実装するためのトレーニングや技術サポートが必要な場合は、営業担当者に連絡するか、このリンク をクリックして、プロフェッショナルサービスの専門家にプロジェクトのカスタム分析を依頼し、見積もりを取得してください。

    専用の Discord チャネル で Data Platform を構築するチームと直接やり取りし、質問をし、フィードバックを提供し、相互作用します。

    OVHcloud サービスのサポートが必要な場合は、ヘルプセンター でリクエストを作成してください。

    ユーザーコミュニティ に参加します。