# SQL Transformations within the OVHcloud Data Platform

* [Introduction](#introduction)
* [Why SQL for Data Transformation?](#why-sql-for-data-transformation)
* [Core Concepts of SQL Transformations](#core-concepts-of-sql-transformations)
    * [Data Cleaning and Preparation](#_1-data-cleaning-and-preparation)
    * [Data Aggregation and Summarization](#_2-data-aggregation-and-summarization)
    * [Data Reshaping and Restructuring](#_3-data-reshaping-and-restructuring)
* [SQL Transformations on OVHcloud Data Platform (Practical Guide)](#sql-transformations-on-ovhcloud-data-platform-practical-guide)
    * [Prerequisites](#prerequisites)
    * [Create a New DPE Notebook](#_1-create-a-new-dpe-notebook)
    * [Connect to the SDK](#_2-connect-to-the-sdk)
    * [List Tables from the Dataset](#_3-list-tables-from-the-dataset)
    * [Connect to the Table and Inspect Data](#_4-connect-to-the-table-and-inspect-data)
    * [Run SQL Commands](#_5-run-sql-commands-transformations)
        * [Simple SQL Query: Daily Total Sales per Location (Initial Exploration)](#_51-simple-sql-query-daily-total-sales-per-location-initial-exploration)
        * [Complex SQL Query: Detailed Cleaned Daily Item Performance](#_52-complex-sql-query-detailed-cleaned-daily-item-performance)
    * [Create a Physical Table from the Transformed Data (CTAS)](#_6-create-a-physical-table-from-the-transformed-data-ctas)
    * [Explore the New Physical Table](#_7-explore-the-new-physical-table)
    * [Create a Logical Object from the Physical Table](#_8-create-a-logical-object-from-the-physical-table-abstraction-layer)

## Introduction

Welcome to the guide on performing **SQL Transformations within the OVHcloud Data Platform**. SQL (Structured Query Language) remains the most widely used and efficient language for querying and transforming structured data. This document will introduce you to the core concepts of using SQL for data manipulation and aggregation directly within our platform, leveraging its powerful capabilities for data processing.

Whether you're looking to clean, reshape, filter, or aggregate large datasets, SQL provides a robust and intuitive way to achieve your data transformation goals. Its declarative nature allows you to focus on what you want to achieve with your data, rather than how the operations are performed, making it accessible to a wide range of data professionals.

## Why SQL for Data Transformation?

SQL is an indispensable tool in the data transformation pipeline for several key reasons:

*   **Universality:** It's a standard language, highly adopted across various databases and data platforms, making skills easily transferable.
*   **Readability & Simplicity:** Its English-like syntax makes complex operations relatively easy to understand and write.
*   **Performance:** SQL engines are highly optimized for relational operations, often outperforming custom code for large-scale data manipulation.
*   **Declarative Nature:** You specify the desired end-state of your data, and the engine determines the most efficient way to achieve it.
*   **Integration:** Seamlessly integrates with data warehousing, business intelligence, and reporting tools.

## Core Concepts of SQL Transformations

At its heart, SQL transformation involves using standard SQL commands to manipulate data. This can include:

### 1. Data Cleaning and Preparation

*   **Filtering:** Using `WHERE` clauses to select specific rows based on conditions.
*   **Selecting/Projecting:** Using `SELECT` statements to choose specific columns and rename them (`AS`).
*   **Type Conversion:** Using `CAST` or `TRY_CAST` functions to change data types. `TRY_CAST` is particularly useful in Trino for gracefully handling conversion errors (returning `NULL` instead of crashing).
*   **Handling Missing Values:** Using `COALESCE` or `CASE` statements to replace `NULL` values.
*   **String Manipulation:** Functions like `SUBSTRING`, `LENGTH`, `UPPER`, `LOWER`, `TRIM` to clean and standardize text data.

### 2. Data Aggregation and Summarization

*   **Grouping:** Using `GROUP BY` to aggregate rows that have the same values in specified columns.
*   **Aggregate Functions:** Applying functions like `COUNT`, `SUM`, `AVG`, `MIN`, `MAX` to summarized groups.
*   **Filtering Aggregations:** Using `HAVING` clauses to filter results of `GROUP BY` operations.

### 3. Data Reshaping and Restructuring

*   **Joins:** Combining data from two or more tables based on related columns (`INNER JOIN`, `LEFT JOIN`, `RIGHT JOIN`, `FULL OUTER JOIN`).
*   **Unions:** Combining the result sets of two or more `SELECT` statements (`UNION`, `UNION ALL`).
*   **Pivoting/Unpivoting:** Transforming rows into columns (pivot) or columns into rows (unpivot) to change the data's structure. In Trino SQL (commonly used in DPE), pivoting is often achieved using `SUM` with `FILTER` or `CASE` statements, and unpivoting with `CROSS JOIN UNNEST`.
*   **Window Functions:** Performing calculations across a set of table rows that are related to the current row, without collapsing rows (`ROW_NUMBER()`, `RANK()`, `LEAD()`, `LAG()`, `SUM() OVER()`, `AVG() OVER()`).

## SQL Transformations on OVHcloud Data Platform (Practical Guide)

This section provides a practical, end-to-end example of performing SQL transformations within the OVHcloud Data Platform's DPE (Data Processing Environment) Notebooks, specifically using the SDK and focusing on the **dirty_cafe_sales** dataset.

### Prerequisites

1. Download the <a href="/en/getting-further/ai-endpoints/resources/dirty_cafe_sales.csv" download>dirty_cafe_sales.csv</a>
2. Upload it to the [Connectors](/en/product/connectors/index.md) and extract the metadata using the [Analyzer](/en/product/connectors/analyzer/index.md).
3. Create a new [Table](/en/product/lakehouse-manager/tables/index.md) from source in the [Lakehouse Manager](/en/product/lakehouse-manager/index.md)

**Dataset Information:**

We will be working with a simulated cafe sales dataset named `dirty_cafe_sales`. It contains the following columns and known data quality issues:

| Column Name      | Description                               | Known Issues                                                              |
| :--------------- | :---------------------------------------- | :------------------------------------------------------------------------ |
| `transaction_id` | Unique identifier for each transaction    | None                                                                      |
| `item`           | Name of the item sold                     | None                                                                      |
| `quantity`       | Number of units sold in the transaction   | Contains 'UNKNOWN' strings                                                |
| `price_per_unit` | Price of a single unit of the item        | None                                                                      |
| `total_spent`    | Total amount spent for the item line in transaction | Contains 'ERROR' strings                                                  |
| `payment_method` | Method used for payment (e.g., Credit Card, Cash) | None                                                                      |
| `location`       | Sales location (e.g., In-store, Takeaway) | Contains 'ERROR' strings and empty strings `''`                             |
| `transaction_date` | Date of the transaction                 | Contains non-date strings (e.g., 'ERROR'), needs casting and error handling |

### 1. Create a New DPE Notebook

*   Navigate to Data Processing Environment (DPE) → Notebooks.
*   Click **+ New Notebook** and for the purpose of this guide we will continue with the **Base Notebook**.
*   Give your notebook a meaningful name (e.g., `Cafe_Sales_SQL_Transformations`).
*   Click on **Create**.
*   Once JupyterLab is open, click on the Python3 Notebook to create a new `.ipynb` file. This is where we are going to run all the following steps in cells.

### 2. Connect to the SDK

We'll use the SDK provided in your DPE environment. This SDK provides methods to connect to the Lakehouse Manager and interact with tables.

```python
# 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. List Tables from the Dataset

Before connecting to `dirty_cafe_sales`, it's good practice to list available tables to confirm its presence and exact name within the specified data path.

```python
# 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. Connect to the Table and Inspect Data

Now, let's connect to the `dirty_cafe_sales` table using `connector.select()` and print its information and descriptive statistics. This step visually confirms the data, including the 'ERROR' and 'UNKNOWN' values we need to clean.

```python
# 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. Run SQL Commands (Transformations)

This is the core of our transformation. We'll define two SQL queries: a simple one for basic exploration and a more complex one for thorough cleaning and detailed aggregation.

> **Important Note on Semicolons:** When executing SQL through an SDK or API in a programmatic environment like a notebook, **do not include a trailing semicolon (`;`)** at the very end of your SQL query string. The API typically expects a single SQL statement without an explicit delimiter at the end. Including it can lead to common errors like `mismatched input ';'` or `syntax error near ';'`.

#### 5.1 Simple SQL Query: Daily Total Sales per Location (Initial Exploration)

This query demonstrates a basic aggregation, initially showing how raw data issues can lead to errors, and then fixing them using `TRY_CAST` for robustness. It also cleans the `location` field.

```python
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 Complex SQL Query: Detailed Cleaned Daily Item Performance

This query will perform robust data cleaning, calculate accurate sales metrics, and aggregate them by date, item, and location. It directly addresses the data quality issues identified (`UNKNOWN` in `quantity`, `ERROR` in `total_spent`, bad `transaction_date` strings, and problematic `location` entries).

```python
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. Create a Physical Table from the Transformed Data (CTAS)

After successfully performing the complex transformation and verifying the results, the next logical step is to persist this cleaned and aggregated data into a new physical table in your database. This is typically done using a `CREATE TABLE AS SELECT` (CTAS) statement. This new table can then be used for reporting, further analysis, or as a source for other data processes, without needing to re-run the complex cleaning logic every time.

> **Important Note on SQL Execution:** Some database connectors or APIs expect only one SQL statement per `query()` call. To execute `DROP TABLE` and `CREATE TABLE AS SELECT`, we'll send them as separate commands. Also, ensure there are no trailing semicolons at the very end of each query string.

```python
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. Explore the New Physical Table

You can use standard SQL metadata commands to explore your newly created physical table. Remember to replace `<your_catalog_name>` and `<your_schema_name>` with your actual values.

```python
# 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. Create a Logical Object from the Physical Table

When you create a new table through SQL, it exists only as a physical table in the database. However, the Tables section in the UI operates at a logical level — it displays tables that are registered as logical objects within the platform.

To make your newly created physical table visible and usable in the UI, you’ll need to create a corresponding logical object. This logical representation acts as a bridge between the database and the platform interface, allowing you to interact with the table’s schema and data directly from the UI.

```python
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}")
```

> if you want to remove a LogicalObject that you have created, you can directly delete the table from the UI or else use the method - `LogicalObject().remove("table_name")` - this will remove the table at both logical and physical levels.

### Conclusion

You have successfully navigated through an end-to-end SQL transformation process on the OVHcloud Data Platform. Starting from a dirty raw dataset, you have applied various SQL cleaning, aggregation, and reshaping techniques to produce a clean, summarized, and highly usable dataset. This transformed data was then persisted into a new physical table and represented as a logical object for seamless integration into your Python workflows. This foundational knowledge empowers you to tackle more complex data preparation challenges and build robust data pipelines.