SQL(構造化問い合わせ言語)は、構造化データの問い合わせと変換に最も広く使用され、効率的な言語です
目的
OVHcloud Data Platform内でのSQL変換のガイドへようこそ。SQL(構造化問い合わせ言語)は、構造化データの問い合わせと変換に最も広く使用され、効率的な言語です。このドキュメントでは、データ操作と集約のためのSQLの基本概念を紹介し、プラットフォーム内でその強力な機能を活用する方法を説明します。
大規模なデータセットをクリーンアップ、再構成、フィルタリング、または集約する必要がある場合、SQLはデータ変換の目標を達成するための堅牢で直感的な方法を提供します。その宣言的な性質により、データ操作の方法ではなく、データで達成したいことを中心に考えることができ、さまざまなデータ専門家にとってアクセスしやすくなります。
SQLは、データ変換パイプラインにおいていくつかの重要な理由から不可欠なツールです:
- 普遍性: 標準言語であり、さまざまなデータベースやデータプラットフォームで広く採用されており、スキルを簡単に移行できます。
- 可読性と簡潔さ: 英語のような構文により、複雑な操作を比較的簡単に理解して記述できます。
- パフォーマンス: SQLエンジンは関係操作に高度に最適化されており、大規模なデータ操作においてカスタムコードを上回ることがよくあります。
- 宣言的な性質: データの最終状態を指定し、エンジンが最も効率的な方法でそれを達成する方法を決定します。
- 統合: データウェアハウス、ビジネスインテリジェンス、レポートツールとシームレスに統合されます。
SQL変換の核心は、標準的なSQLコマンドを使用してデータを操作することです。これは次のようなものを含むことがあります。
1. データのクリーンアップと準備
- フィルタリング:
WHERE 句を使用して、条件に基づいて特定の行を選択します。
- 選択/投影:
SELECT 文を使用して特定の列を選択し、名前を変更します(AS)。
- 型変換:
CAST または TRY_CAST 関数を使用してデータ型を変更します。TRY_CAST は、Trinoで変換エラーを優雅に処理する(NULL を返す代わりにクラッシュする)ために特に便利です。
- 欠損値の処理:
COALESCE または CASE 文を使用して NULL 値を置き換えます。
- 文字列操作:
SUBSTRING、LENGTH、UPPER、LOWER、TRIM のような関数を使用してテキストデータをクリーンアップし、標準化します。
2. データの集約と要約
- グループ化:
GROUP BY を使用して、指定された列で同じ値を持つ行を集約します。
- 集約関数:
COUNT、SUM、AVG、MIN、MAX のような関数を適用して集約されたグループを要約します。
- 集約のフィルタリング:
HAVING 句を使用して GROUP BY 操作の結果をフィルタリングします。
3. データの再構成と再構造化
- 結合: 関連する列に基づいて2つ以上のテーブルのデータを結合します(
INNER JOIN、LEFT JOIN、RIGHT JOIN、FULL OUTER JOIN)。
- ユニオン: 2つ以上の
SELECT 文の結果セットを結合します(UNION、UNION ALL)。
- ピボット/アンピボット: 行を列(ピボット)または列を行(アンピボット)に変換してデータの構造を変更します。Trino SQL(DPEで一般的に使用される)では、ピボットは
SUM と FILTER または CASE 文を使用してよく実現され、アンピボットは CROSS JOIN UNNEST を使用して実現されます。
- ウィンドウ関数: 現在の行に関連するテーブル行のセット全体で計算を実行し、行を折りたたむことなく(
ROW_NUMBER()、RANK()、LEAD()、LAG()、SUM() OVER()、AVG() OVER())。
このセクションでは、OVHcloud Data Platform の DPE(Data Processing Environment)ノートブック内で SQL 変換を実行する、実践的なエンドツーエンドの例を提供します。具体的には、SDK を使用し、dirty_cafe_sales データセットに焦点を当てます。
前提条件
- dirty_cafe_sales.csv をダウンロードします。
- Connectors にアップロードし、Analyzer を使用してメタデータを抽出します。
- Lakehouse Manager で Table をソースから新規作成します。
データセット情報:
私たちは、dirty_cafe_sales という名前のシミュレートされたカフェ販売データセットを使用します。このデータセットには、次の列と既知のデータ品質の問題が含まれています。
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())
これが変換の核心です。基本的な探索用のシンプルな 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)
このクエリは、データを強固にクリーンアップし、正確な販売メトリクスを計算し、日付、アイテム、場所別に集計します。これは、quantity の UNKNOWN、total_spent の ERROR、悪い 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)
複雑な変換を成功裏に実行し、結果を確認した後、次の論理的なステップは、データベース内の新しい物理テーブルにこのクリーンアップ済みの集計データを保持することです。これは、通常 CREATE TABLE AS SELECT(CTAS)ステートメントを使用して実行されます。この新しいテーブルは、レポート作成、さらに分析、または他のデータ処理のソースとして使用できます。複雑なクリーンアップロジックを毎回再実行する必要はありません。
Info
SQL 実行に関する重要な注意点: 一部のデータベースコネクタまたは API は、query() 呼び出しごとに 1 つの SQL ステートメントのみを期待します。DROP TABLE と CREATE 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 サービスのサポートが必要な場合は、ヘルプセンター でリクエストを作成してください。
ユーザーコミュニティ に参加します。