Trong phần này, chúng ta sẽ xây dựng ETL pipeline để chuyển đổi dữ liệu clickstream thô thành chuỗi thời gian có cấu trúc cho machine learning. Pipeline sẽ thực hiện data cleaning, aggregation, feature engineering và output dữ liệu dưới dạng Parquet partitioned để tối ưu performance. Đây là bước quan trọng để chuẩn bị dữ liệu phù hợp cho việc training forecasting models.
✅ Chuyển clickstream thành chuỗi thời gian ngày × KEY
✅ Tạo features và labels phù hợp cho forecasting
✅ Thực hiện data quality checks và validation
✅ Output Parquet partitioned cho performance tối ưu
Input Schema (Raw Clickstream):
Output Schema (Daily Aggregated):
1. Target Labels:
count(event_type == 'purchase') theo ngày × KEY2. Exogenous Variables:
3. Conversion Metrics:
purchases_t / views_tcarts_t / views_t4. Lag Features:
5. Moving Averages:
6. Calendar Features:
Step 1: Data Cleaning
Step 2: Aggregation
Step 3: Feature Engineering
Step 4: Quality Checks
Step 5: Output
Required Components:
retail_data:1 (từ Task 4: Data Asset Registration)cpu-cluster (từ Task 6)retail-train-env với pandas/dask (từ Task 6)Basic Configuration:
etl_aggregateretail-train-env (pandas/dask)cpu-clusterETL Aggregate PipelineInputs:
retail_data:1product_id|category_code|brandOutputs:
aggregates/ in datastoreretail_aggregates:1 (auto-registered)Command:
python etl_aggregate.py --input_data ${{inputs.retail_data}} --key_type ${{inputs.key_type}} --output_path ${{outputs.aggregates}}
Step 1: Navigate to Jobs
Azure ML Studio → Jobs → +New → Command job
Step 2: Basic Configuration
etl_aggregateretail-train-envcpu-clusterETL Aggregate PipelineStep 3: Inputs Configuration
retail_data:1product_id, category_code, brandStep 4: Outputs Configuration
aggregates/retail_aggregates:1Step 5: Command Configuration
python etl_aggregate.py \
--input_data ${{inputs.retail_data}} \
--key_type ${{inputs.key_type}} \
--output_path ${{outputs.aggregates}}
# etl_aggregate.py
import pandas as pd
import dask.dataframe as dd
import numpy as np
import argparse
from datetime import datetime, timedelta
import logging
from pathlib import Path
def setup_logging():
"""Setup logging configuration"""
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(levelname)s - %(message)s'
)
return logging.getLogger(__name__)
def load_and_clean_data(input_path, logger):
"""Load and clean raw clickstream data"""
logger.info("Loading raw clickstream data...")
# Load data with Dask for large datasets
df = dd.read_parquet(f"{input_path}/**/*.parquet")
logger.info(f"Loaded {len(df)} records")
# Data cleaning
logger.info("Performing data cleaning...")
# Convert to pandas for cleaning operations
df = df.compute()
# Convert timezone to UTC
df['event_time'] = pd.to_datetime(df['event_time']).dt.tz_localize('UTC')
# Remove invalid records
initial_count = len(df)
# Filter out invalid prices
df = df[df['price'] > 0]
# Filter out empty IDs
df = df[df['product_id'].notna()]
df = df[df['product_id'] != '']
# Filter valid event types
valid_events = ['view', 'cart', 'purchase']
df = df[df['event_type'].isin(valid_events)]
# Extract date
df['date'] = df['event_time'].dt.date
final_count = len(df)
dropped_pct = (initial_count - final_count) / initial_count * 100
logger.info(f"Dropped {initial_count - final_count} records ({dropped_pct:.2f}%)")
logger.info(f"Final dataset: {final_count} records")
return df
def aggregate_by_key(df, key_type, logger):
"""Aggregate data by date and key"""
logger.info(f"Aggregating by {key_type}...")
# Group by date and key
grouped = df.groupby(['date', key_type]).agg({
'event_type': [
('purchases', lambda x: (x == 'purchase').sum()),
('views', lambda x: (x == 'view').sum()),
('carts', lambda x: (x == 'cart').sum())
],
'price': 'mean'
}).reset_index()
# Flatten column names
grouped.columns = ['date', 'key_value', 'purchases_t', 'views_t', 'carts_t', 'price_mean_t']
# Add key type
grouped['key_type'] = key_type
# Calculate conversion rates
grouped['conv_rate'] = grouped['purchases_t'] / grouped['views_t'].replace(0, np.nan)
grouped['cart_rate'] = grouped['carts_t'] / grouped['views_t'].replace(0, np.nan)
# Sort by date and key
grouped = grouped.sort_values(['date', 'key_value'])
logger.info(f"Aggregated to {len(grouped)} records")
logger.info(f"Date range: {grouped['date'].min()} to {grouped['date'].max()}")
logger.info(f"Unique keys: {grouped['key_value'].nunique()}")
return grouped
def add_features(df, logger):
"""Add lag features and moving averages"""
logger.info("Adding engineered features...")
# Sort by key_value and date
df = df.sort_values(['key_value', 'date']).reset_index(drop=True)
# Add lag features
df['purchases_t-1'] = df.groupby('key_value')['purchases_t'].shift(1)
df['purchases_t-7'] = df.groupby('key_value')['purchases_t'].shift(7)
df['views_t-1'] = df.groupby('key_value')['views_t'].shift(1)
# Add moving averages
df['MA7'] = df.groupby('key_value')['purchases_t'].rolling(window=7, min_periods=1).mean().reset_index(0, drop=True)
df['MA28'] = df.groupby('key_value')['purchases_t'].rolling(window=28, min_periods=1).mean().reset_index(0, drop=True)
# Add calendar features
df['date'] = pd.to_datetime(df['date'])
df['day_of_week'] = df['date'].dt.dayofweek
df['is_weekend'] = df['day_of_week'].isin([5, 6])
# Add holiday flag (simplified - can be enhanced with external calendar)
df['holiday_flag'] = 0 # Placeholder for holiday logic
logger.info("Feature engineering completed")
return df
def quality_checks(df, logger):
"""Perform data quality checks"""
logger.info("Performing quality checks...")
# Log statistics
logger.info(f"Dataset statistics:")
logger.info(f" Total records: {len(df)}")
logger.info(f" Date range: {df['date'].min()} to {df['date'].max()}")
logger.info(f" Unique keys: {df['key_value'].nunique()}")
logger.info(f" Key type: {df['key_type'].iloc[0]}")
# Check for missing values
missing_counts = df.isnull().sum()
if missing_counts.sum() > 0:
logger.warning(f"Missing values found:")
for col, count in missing_counts[missing_counts > 0].items():
logger.warning(f" {col}: {count}")
# Check data consistency
negative_values = (df[['purchases_t', 'views_t', 'carts_t']] < 0).sum().sum()
if negative_values > 0:
logger.error(f"Found {negative_values} negative values in count columns")
logger.info("Quality checks completed")
def save_output(df, output_path, logger):
"""Save aggregated data as partitioned Parquet files"""
logger.info(f"Saving output to {output_path}...")
# Create output directory
Path(output_path).mkdir(parents=True, exist_ok=True)
# Convert date back to string for partitioning
df['date_str'] = df['date'].dt.strftime('%Y-%m-%d')
# Save as partitioned Parquet
df.to_parquet(
output_path,
partition_cols=['key_type', 'date_str'],
index=False
)
logger.info("Output saved successfully")
# List output files
output_files = list(Path(output_path).rglob('*.parquet'))
logger.info(f"Created {len(output_files)} Parquet files")
def main():
parser = argparse.ArgumentParser(description='ETL Aggregate Pipeline')
parser.add_argument('--input_data', required=True, help='Input data path')
parser.add_argument('--key_type', required=True, choices=['product_id', 'category_code', 'brand'], help='Key type for aggregation')
parser.add_argument('--output_path', required=True, help='Output path')
args = parser.parse_args()
# Setup logging
logger = setup_logging()
try:
# Load and clean data
df = load_and_clean_data(args.input_data, logger)
# Aggregate by key
df_agg = aggregate_by_key(df, args.key_type, logger)
# Add features
df_features = add_features(df_agg, logger)
# Quality checks
quality_checks(df_features, logger)
# Save output
save_output(df_features, args.output_path, logger)
logger.info("ETL pipeline completed successfully!")
except Exception as e:
logger.error(f"ETL pipeline failed: {str(e)}")
raise
if __name__ == "__main__":
main()
# retail-train-env.yml
name: retail-train-env
channels:
- conda-forge
- defaults
dependencies:
- python=3.8
- pandas>=1.3.0
- dask>=2021.0.0
- numpy>=1.21.0
- pyarrow>=5.0.0
- fastparquet>=0.7.0
- pip
- pip:
- azure-ai-ml>=1.0.0
- azure-identity>=1.7.0
# Submit job via Azure CLI
az ml job create \
--file etl_aggregate_job.yml \
--resource-group retail-dev-rg \
--workspace-name retail-ml-workspace
# etl_aggregate_job.yml
$schema: https://azuremlschemas.azureedge.net/latest/commandJob.schema.json
command: >
python etl_aggregate.py
--input_data ${{inputs.retail_data}}
--key_type ${{inputs.key_type}}
--output_path ${{outputs.aggregates}}
code: ./src
environment: azureml://registries/azureml/environments/retail-train-env/labels/latest
compute: azureml://subscriptions/{subscription-id}/resourceGroups/retail-dev-rg/providers/Microsoft.MachineLearningServices/workspaces/retail-ml-workspace/computes/cpu-cluster
inputs:
retail_data:
type: uri_folder
path: azureml://datastores/workspaceblobstore/paths/raw/ecommerce/
key_type:
type: string
default: product_id
outputs:
aggregates:
type: uri_folder
display_name: ETL Aggregate Pipeline
description: Transform clickstream data to time series format
tags:
task: etl
pipeline: aggregate
dataset: retail
# monitor_job.py
from azure.ai.ml import MLClient
from azure.identity import DefaultAzureCredential
import time
def monitor_job(job_name):
credential = DefaultAzureCredential()
ml_client = MLClient(
credential=credential,
subscription_id="your-subscription-id",
resource_group_name="retail-dev-rg",
workspace_name="retail-ml-workspace"
)
job = ml_client.jobs.get(job_name)
print(f"Job Status: {job.status}")
print(f"Start Time: {job.creation_context.created_at}")
while job.status in ['NotStarted', 'Starting', 'Preparing', 'Running']:
time.sleep(30)
job = ml_client.jobs.get(job_name)
print(f"Status: {job.status}")
if job.status == 'Failed':
print("Job failed!")
break
elif job.status == 'Completed':
print("Job completed successfully!")
break
return job
if __name__ == "__main__":
job = monitor_job("etl_aggregate")
etl_aggregateretail-train-envcpu-clusterretail_data:1etl_aggregateCompletedretail_aggregates:1aggregates/etl_aggregate Run Succeededretail_aggregates:1 registered#!/bin/bash
# submit-etl-job.sh
# Configuration
SUBSCRIPTION_ID="your-subscription-id"
RESOURCE_GROUP="retail-dev-rg"
WORKSPACE_NAME="retail-ml-workspace"
JOB_NAME="etl_aggregate"
echo "🚀 Submitting ETL Aggregate Job..."
# Login to Azure
az login
# Set subscription
az account set --subscription $SUBSCRIPTION_ID
# Submit job
az ml job create \
--file etl_aggregate_job.yml \
--resource-group $RESOURCE_GROUP \
--workspace-name $WORKSPACE_NAME
# Monitor job
echo "📊 Monitoring job execution..."
az ml job show --name $JOB_NAME --resource-group $RESOURCE_GROUP --workspace-name $WORKSPACE_NAME
echo "✅ Job submitted successfully!"
# manage_etl_jobs.py
from azure.ai.ml import MLClient, command
from azure.ai.ml.entities import Data
from azure.identity import DefaultAzureCredential
def submit_etl_job(key_type='product_id'):
credential = DefaultAzureCredential()
ml_client = MLClient(
credential=credential,
subscription_id="your-subscription-id",
resource_group_name="retail-dev-rg",
workspace_name="retail-ml-workspace"
)
# Create command job
job = command(
code="./src",
command="python etl_aggregate.py --input_data ${{inputs.retail_data}} --key_type ${{inputs.key_type}} --output_path ${{outputs.aggregates}}",
inputs={
"retail_data": Data(type="uri_folder", path="azureml://datastores/workspaceblobstore/paths/raw/ecommerce/"),
"key_type": key_type
},
outputs={
"aggregates": {"type": "uri_folder"}
},
environment="azureml://registries/azureml/environments/retail-train-env/labels/latest",
compute="cpu-cluster",
display_name=f"ETL Aggregate - {key_type}"
)
# Submit job
submitted_job = ml_client.jobs.create_or_update(job)
print(f"✅ Job submitted: {submitted_job.name}")
return submitted_job
def submit_all_key_types():
"""Submit ETL jobs for all key types"""
key_types = ['product_id', 'category_code', 'brand']
for key_type in key_types:
job = submit_etl_job(key_type)
print(f"Submitted job for {key_type}: {job.name}")
if __name__ == "__main__":
submit_all_key_types()
Issue 1: “Job failed with memory error”
# Increase compute instance size
# Use Dask for large datasets
# Implement data chunking
Issue 2: “Data Asset not found”
# Verify data asset exists
az ml data show --name retail_data --version 1
# Check workspace context
az ml workspace show --name retail-ml-workspace --resource-group retail-dev-rg
Issue 3: “Output path not accessible”
# Check compute permissions
# Verify datastore configuration
# Ensure output folder exists
Best Practice: Always validate ETL output data quality before proceeding to model training. Use the quality checks and statistics to ensure data integrity.
Performance Note: For large datasets, consider using Dask or Spark for distributed processing. Monitor memory usage and adjust compute resources accordingly.
ETL Aggregate pipeline hoàn tất! 🎉 Time series dataset đã sẵn sàng cho Task 6: Model Training.