Skip to content
IRC-CodingIRC-Coding
Fundamentos Data ScienceBig Data AnalyticsML PipelinesProcesos ETLData VisualizationAlgoritmosAlgoritmoFundamentos

Fundamentos de Data Science: Analytics, ML y ETL

Aprende Data Science con Big Data Analytics, ML Pipelines, ETL y visualización de datos. Ejemplos prácticos con Python, Pandas y Scikit-learn.

S

schutzgeist

23 min read
Fundamentos de Data Science: Analytics, ML y ETL

Fundamentos de Data Science: Big Data Analytics, ML-Pipelines, ETL y Visualización

Este artículo es una introducción exhaustiva a los fundamentos de Data Science, incluyendo Big Data Analytics, Machine Learning Pipelines, procesos ETL y Data Visualization con ejemplos prácticos.

En Resumen

Data Science combina estadística, programación y conocimiento del dominio para extraer información valiosa de los datos. Big Data procesa volúmenes masivos, ML-Pipelines automatizan modelos, ETL prepara los datos, Visualization hace visibles los resultados.

Descripción Técnica Compacta

Data Science es un campo interdisciplinario que utiliza métodos científicos, procesos, algoritmos y sistemas para extraer conocimientos e información de datos estructurados y no estructurados.

Áreas Clave:

Big Data Analytics

  • Concepto: Procesamiento de volúmenes extremadamente grandes de datos
  • Modelo 3V: Volume, Velocity, Variety
  • Tecnologías: Hadoop, Spark, bases de datos NoSQL
  • Aplicaciones: Análisis en tiempo real, Predictive Analytics

Machine Learning Pipelines

  • Concepto: Orquestación automatizada de flujos de trabajo ML
  • Fases: Data Collection → Preprocessing → Training → Evaluation → Deployment
  • Herramientas: Scikit-learn, TensorFlow, MLflow, Airflow
  • MLOps: Versionamiento, Monitoring, Retraining

Procesos ETL

  • Extract: Extraer datos de diversas fuentes
  • Transform: Limpiar y transformar datos
  • Load: Cargar datos en sistemas destino
  • Herramientas: Apache NiFi, Talend, AWS Glue

Data Visualization

  • Concepto: Representación visual de datos e información
  • Tipos de Gráficos: Bar, Line, Scatter, Heatmap, Treemap
  • Herramientas: Matplotlib, Seaborn, Plotly, Tableau
  • Principios: Claridad, Precisión, Eficiencia

Puntos Clave para Evaluaciones

  • Data Science: Campo interdisciplinario para análisis de datos
  • Big Data: Volúmenes grandes, rápidos y variados de datos
  • Machine Learning: Reconocimiento automatizado de patrones en datos
  • ETL: Extract-Transform-Load para preparación de datos
  • Visualization: Representación visual de datos
  • Pipelines: Procesos automatizados de procesamiento de datos
  • Analytics: Análisis de datos estadístico y exploratorio
  • Relevancia IHK: Análisis de datos moderno e inteligencia empresarial

Componentes Principales

  1. Data Collection: Fuentes y recopilación de datos
  2. Data Storage: Almacenamiento y organización
  3. Data Processing: Limpieza y transformación
  4. Data Analysis: Análisis estadístico y exploratorio
  5. Machine Learning: Modelos y algoritmos
  6. Data Visualization: Representación visual
  7. Model Deployment: Despliegue en producción
  8. Monitoring: Detección de rendimiento y drift

Ejemplos Prácticos

1. Procesamiento de Big Data con Apache Spark y Python

from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.types import *
from pyspark.ml.feature import VectorAssembler, StandardScaler
from pyspark.ml.regression import LinearRegression
from pyspark.ml.evaluation import RegressionEvaluator
import matplotlib.pyplot as plt
import seaborn as sns

# Big Data Processing Demo
class BigDataAnalytics:
    
    def __init__(self, app_name="BigDataAnalytics"):
        """Inicializar Spark Session"""
        self.spark = SparkSession.builder \
            .appName(app_name) \
            .config("spark.sql.adaptive.enabled", "true") \
            .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
            .config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB") \
            .getOrCreate()
        
        print(f"Spark Session creada: {self.spark.version}")
    
    def create_sample_data(self):
        """Crear datos de ejemplo para demostración de Big Data"""
        import random
        from datetime import datetime, timedelta
        
        # Datos de ventas simulados
        data = []
        products = ["Laptop", "Smartphone", "Tablet", "Headphones", "Mouse", "Keyboard"]
        regions = ["North", "South", "East", "West", "Central"]
        
        start_date = datetime(2023, 1, 1)
        
        for i in range(1000000):  # 1M records
            date = start_date + timedelta(days=random.randint(0, 365))
            record = {
                "transaction_id": f"T{i:06d}",
                "date": date.strftime("%Y-%m-%d"),
                "product": random.choice(products),
                "region": random.choice(regions),
                "quantity": random.randint(1, 10),
                "unit_price": round(random.uniform(10, 1000), 2),
                "customer_id": f"C{random.randint(1, 50000):05d}",
                "sales_rep": f"SR{random.randint(1, 100):03d}"
            }
            record["total_amount"] = record["quantity"] * record["unit_price"]
            data.append(record)
        
        # Crear DataFrame de Spark
        schema = StructType([
            StructField("transaction_id", StringType(), True),
            StructField("date", StringType(), True),
            StructField("product", StringType(), True),
            StructField("region", StringType(), True),
            StructField("quantity", IntegerType(), True),
            StructField("unit_price", DoubleType(), True),
            StructField("customer_id", StringType(), True),
            StructField("sales_rep", StringType(), True),
            StructField("total_amount", DoubleType(), True)
        ])
        
        df = self.spark.createDataFrame(data, schema)
        print(f"Datos de ejemplo creados: {df.count()} filas")
        
        return df
    
    def basic_analytics(self, df):
        """Realizar análisis básicos"""
        print("=== Análisis Básicos ===")
        
        # Mostrar esquema
        print("Esquema:")
        df.printSchema()
        
        # Resumen estadístico
        print("\nResumen Estadístico:")
        df.describe().show()
        
        # Top productos por ingresos
        print("\nTop Productos por Ingresos:")
        product_sales = df.groupBy("product") \
            .agg(sum("total_amount").alias("total_sales"),
                 count("transaction_id").alias("transaction_count")) \
            .orderBy(col("total_sales").desc())
        
        product_sales.show()
        
        # Análisis regional
        print("\nAnálisis de Ingresos Regionales:")
        regional_sales = df.groupBy("region") \
            .agg(sum("total_amount").alias("total_sales"),
                 avg("total_amount").alias("avg_transaction"),
                 count("transaction_id").alias("transaction_count")) \
            .orderBy(col("total_sales").desc())
        
        regional_sales.show()
        
        # Análisis temporal
        print("\nEvolución de Ingresos Mensuales:")
        monthly_sales = df.withColumn("month", substring(col("date"), 1, 7)) \
            .groupBy("month") \
            .agg(sum("total_amount").alias("monthly_sales"),
                 count("transaction_id").alias("transaction_count")) \
            .orderBy("month")
        
        monthly_sales.show()
        
        return product_sales, regional_sales, monthly_sales
    
    def advanced_analytics(self, df):
        """Analytics avanzados con Window Functions"""
        print("\n=== Analytics Avanzados ===")
        
        # Window Functions para sumas acumulativas
        from pyspark.sql.window import Window
        
        # Ingresos mensuales acumulativos
        window_spec = Window.partitionBy("month").orderBy("date") \
            .rowsBetween(Window.unboundedPreceding, Window.currentRow)
        
        monthly_cumulative = df.withColumn("month", substring(col("date"), 1, 7)) \
            .groupBy("month", "date") \
            .agg(sum("total_amount").alias("daily_sales")) \
            .withColumn("cumulative_sales", sum("daily_sales").over(window_spec)) \
            .orderBy("month", "date")
        
        print("Ingresos Mensuales Acumulativos:")
        monthly_cumulative.show(20)
        
        # Sales Reps con Mejor Desempeño
        rep_performance = df.groupBy("sales_rep") \
            .agg(sum("total_amount").alias("total_sales"),
                 count("transaction_id").alias("transaction_count"),
                 avg("total_amount").alias("avg_transaction")) \
            .withColumn("performance_rank", percent_rank().over(Window.orderBy(col("total_sales").desc()))) \
            .filter(col("performance_rank") <= 0.1)  # Top 10%
        
        print("\nTop 10% Sales Reps:")
        rep_performance.show()
        
        # Correlaciones de Productos
        product_correlation = df.groupBy("customer_id", "product") \
            .agg(sum("quantity").alias("product_quantity")) \
            .groupBy("customer_id") \
            .pivot("product") \
            .agg(sum("product_quantity")) \
            .fillna(0)
        
        print("\nMatriz de Correlación de Productos (Top 10 Clientes):")
        product_correlation.limit(10).show()
        
        return monthly_cumulative, rep_performance, product_correlation
    
    def machine_learning_pipeline(self, df):
        """Machine Learning Pipeline para predicción de ingresos"""
        print("\n=== Machine Learning Pipeline ===")
        
        # Preparar features para ML
        feature_data = df.groupBy("product", "region", "month") \
            .agg(sum("total_amount").alias("target_sales"),
                 count("transaction_id").alias("transaction_count"),
                 avg("unit_price").alias("avg_unit_price"),
                 sum("quantity").alias("total_quantity"))
        
        # Codificar features categóricas
        from pyspark.ml.feature import StringIndexer, OneHotEncoder
        
        # Product Indexer
        product_indexer = StringIndexer(inputCol="product", outputCol="product_index")
        product_indexed = product_indexer.fit(feature_data).transform(feature_data)
        
        # Region Indexer
        region_indexer = StringIndexer(inputCol="region", outputCol="region_index")
        region_indexed = region_indexer.fit(product_indexed).transform(product_indexed)
        
        # One-Hot Encoding
        encoder = OneHotEncoder(inputCols=["product_index", "region_index"],
                               outputCols=["product_encoded", "region_encoded"])
        encoded_data = encoder.fit(region_indexed).transform(region_indexed)
        
        # Crear Vector de Features
        assembler = VectorAssembler(
            inputCols=["transaction_count", "avg_unit_price", "total_quantity",
                       "product_encoded", "region_encoded"],
            outputCol="features"
        )
        
        assembled_data = assembler.transform(encoded_data)
        
        # Escalar features
        scaler = StandardScaler(inputCol="features", outputCol="scaled_features")
        scaler_model = scaler.fit(assembled_data)
        scaled_data = scaler_model.transform(assembled_data)
        
        # Train-Test Split
        train_data, test_data = scaled_data.randomSplit([0.8, 0.2], seed=42)
        
        print(f"Datos de entrenamiento: {train_data.count()} filas")
        print(f"Datos de prueba: {test_data.count()} filas")
        
        # Regresión Lineal
        lr = LinearRegression(featuresCol="scaled_features", 
                            labelCol="target_sales",
                            predictionCol="predicted_sales")
        
        lr_model = lr.fit(train_data)
        
        # Predicciones en datos de prueba
        predictions = lr_model.transform(test_data)
        
        # Evaluar modelo
        evaluator = RegressionEvaluator(labelCol="target_sales",
                                     predictionCol="predicted_sales",
                                     metricName="rmse")
        
        rmse = evaluator.evaluate(predictions)
        r2 = evaluator.setMetricName("r2").evaluate(predictions)
        
        print(f"Evaluación del Modelo:")
        print(f"RMSE: {rmse:.2f}")
        print(f"R²: {r2:.4f}")
        
        # Mostrar Feature Importance
        print(f"\nCoeficientes de Features:")
        feature_names = ["transaction_count", "avg_unit_price", "total_quantity",
                       "product_encoded", "region_encoded"]
        
        for i, (name, coef) in enumerate(zip(feature_names, lr_model.coefficients)):
            print(f"{name}: {coef:.4f}")
        
        return lr_model, predictions, scaler_model
    
    def real_time_processing_simulation(self, df):
        """Simular procesamiento en tiempo real"""
        print("\n=== Simulación de Procesamiento en Tiempo Real ===")
        
        # Simular DataFrame streaming
        # En la práctica, esto vendría de Kafka, Kinesis, etc.
        
        # Agregaciones para monitoreo en tiempo real
        real_time_metrics = df.groupBy("product", "region") \
            .agg(count("transaction_id").alias("transaction_count"),
                 sum("total_amount").alias("total_sales"),
                 avg("total_amount").alias("avg_transaction")) \
            .orderBy(col("total_sales").desc())
        
        print("Métricas en Tiempo Real:")
        real_time_metrics.show()
        
        # Detección de Anomalías (simplificada)
        from pyspark.sql.functions import stddev
        
        # Límites estadísticos para detección de anomalías
        stats = df.groupBy("product") \
            .agg(avg("total_amount").alias("avg_amount"),
                 stddev("total_amount").alias("std_amount")) \
            .withColumn("upper_bound", col("avg_amount") + 2 * col("std_amount")) \
            .withColumn("lower_bound", col("avg_amount") - 2 * col("std_amount"))
        
        print("Límites de Anomalías por Producto:")
        stats.show()
        
        return real_time_metrics, stats
    
    def data_export(self, df, predictions):
        """Exportar datos para procesamiento posterior"""
        print("\n=== Exportar Datos ===")
        
        # Exportar resultados en diferentes formatos
        output_path = "/tmp/big_data_results"
        
        # Resultados de Analytics
        product_sales = df.groupBy("product") \
            .agg(sum("total_amount").alias("total_sales")) \
            .orderBy(col("total_sales").desc())
        
        product_sales.write.mode("overwrite").parquet(f"{output_path}/product_sales")
        print(f"Ingresos por Producto exportados a: {output_path}/product_sales")
        
        # Predicciones ML
        predictions.select("product", "region", "target_sales", "predicted_sales") \
            .write.mode("overwrite").parquet(f"{output_path}/predictions")
        
        print(f"Predicciones ML exportadas a: {output_path}/predictions")
        
        # Estadísticas Resumidas
        summary_stats = df.agg(
            count("transaction_id").alias("total_transactions"),
            sum("total_amount").alias("total_revenue"),
            avg("total_amount").alias("avg_transaction_value"),
            countDistinct("customer_id").alias("unique_customers"),
            countDistinct("product").alias("unique_products")
        )
        
        summary_stats.write.mode("overwrite").json(f"{output_path}/summary_stats")
        print(f"Estadísticas Resumidas exportadas a: {output_path}/summary_stats")
    
    def cleanup(self):
        """Liberar recursos"""
        self.spark.stop()
        print("Spark Session finalizada")

# Ejecutar demostración de Big Data
def big_data_demo():
    analytics = BigDataAnalytics()
    
    try:
        # Crear datos
        df = analytics.create_sample_data()
        
        # Analytics básicos
        product_sales, regional_sales, monthly_sales = analytics.basic_analytics(df)
        
        # Analytics avanzados
        monthly_cumulative, rep_performance, product_correlation = analytics.advanced_analytics(df)
        
        # Machine Learning Pipeline
        lr_model, predictions, scaler_model = analytics.machine_learning_pipeline(df)
        
        # Procesamiento en tiempo real
        real_time_metrics, stats = analytics.real_time_processing_simulation(df)
        
        # Exportar datos
        analytics.data_export(df, predictions)
        
        print("\nDemostración de Big Data Analytics completada exitosamente!")
        
    except Exception as e:
        print(f"Error en demostración de Big Data: {e}")
    finally:
        analytics.cleanup()

if __name__ == "__main__":
    big_data_demo()

2. Pipeline ETL con Apache Airflow

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow.providers.postgres.hooks.postgres import PostgresHook
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
from airflow.providers.amazon.aws.operators.s3 import S3UploadOperator
from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator
from airflow.providers.google.cloud.hooks.bigquery import BigQueryHook
from airflow.models import Variable
from airflow.utils.dates import days_ago
import pandas as pd
import numpy as np
from datetime import datetime, timedelta
import json
import logging

# Configuración del Pipeline ETL
default_args = {
    'owner': 'data-science-team',
    'depends_on_past': False,
    'start_date': days_ago(1),
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

# Definición del DAG
dag = DAG(
    'data_science_etl_pipeline',
    default_args=default_args,
    description='ETL Pipeline für Data Science Analytics',
    schedule_interval='@daily',
    catchup=False,
    tags=['data-science', 'etl', 'analytics'],
)

# Configurar logging
logger = logging.getLogger(__name__)

def extract_customer_data(**context):
    """Extrae datos de clientes desde varias fuentes"""
    logger.info("Extrayendo datos de clientes...")
    
    # PostgreSQL Hook
    postgres_hook = PostgresHook(postgres_conn_id='postgres_default')
    
    # Extraer datos de clientes desde PostgreSQL
    customer_query = """
        SELECT 
            customer_id,
            first_name,
            last_name,
            email,
            phone,
            registration_date,
            last_login_date,
            total_orders,
            total_spent,
            customer_segment
        FROM customers 
        WHERE last_login_date >= DATE_SUB(CURRENT_DATE, INTERVAL 30 DAY)
    """
    
    customers_df = postgres_hook.get_pandas_df(customer_query)
    logger.info(f"{len(customers_df)} registros de clientes extraídos")
    
    # S3 Hook para datos adicionales
    s3_hook = S3Hook(aws_conn_id='aws_default')
    
    # Datos de interacciones de productos desde S3
    bucket_name = Variable.get('data_bucket_name', default_var='analytics-data')
    
    try:
        # Cargar archivo CSV desde S3
        file_content = s3_hook.read_key(key='customer_interactions.csv', bucket_name=bucket_name)
        interactions_df = pd.read_csv(pd.StringIO(file_content))
        logger.info(f"{len(interactions_df)} registros de interacciones cargados desde S3")
    except Exception as e:
        logger.warning(f"No se encontraron datos de interacciones en S3: {e}")
        interactions_df = pd.DataFrame()
    
    # Combinar datos
    if not interactions_df.empty:
        combined_df = customers_df.merge(
            interactions_df, 
            on='customer_id', 
            how='left'
        )
    else:
        combined_df = customers_df
    
    # Guardar datos como JSON para la siguiente tarea
    combined_json = combined_df.to_json(date_format='iso')
    
    # Guardar en XCom
    task_instance = context['task_instance']
    task_instance.xcom_push(key='customer_data', value=combined_json)
    
    logger.info("Extracción de datos de clientes completada")
    return combined_json

def transform_customer_data(**context):
    """Transforma y limpia los datos de clientes"""
    logger.info("Transformando datos de clientes...")
    
    # Obtener datos de la tarea anterior
    task_instance = context['task_instance']
    customer_json = task_instance.xcom_pull(task_ids='extract_customer_data', key='customer_data')
    
    # JSON a DataFrame
    customers_df = pd.read_json(customer_json)
    
    # Limpieza de datos
    logger.info("Realizando limpieza de datos...")
    
    # 1. Eliminar duplicados
    original_count = len(customers_df)
    customers_df = customers_df.drop_duplicates(subset=['customer_id'])
    duplicates_removed = original_count - len(customers_df)
    logger.info(f"{duplicates_removed} duplicados eliminados")
    
    # 2. Tratar valores faltantes
    # Columnas numéricas: imputar con mediana
    numeric_columns = ['total_orders', 'total_spent']
    for col in numeric_columns:
        if col in customers_df.columns:
            median_value = customers_df[col].median()
            customers_df[col].fillna(median_value, inplace=True)
    
    # Columnas categóricas: imputar con moda
    categorical_columns = ['customer_segment']
    for col in categorical_columns:
        if col in customers_df.columns:
            mode_value = customers_df[col].mode()[0] if not customers_df[col].mode().empty else 'Unknown'
            customers_df[col].fillna(mode_value, inplace=True)
    
    # 3. Normalizar columnas de fechas
    date_columns = ['registration_date', 'last_login_date']
    for col in date_columns:
        if col in customers_df.columns:
            customers_df[col] = pd.to_datetime(customers_df[col], errors='coerce')
    
    # 4. Ingeniería de características
    logger.info("Creando nuevas características...")
    
    # Calcular Customer Lifetime Value (CLV)
    if 'total_spent' in customers_df.columns and 'total_orders' in customers_df.columns:
        customers_df['avg_order_value'] = customers_df['total_spent'] / customers_df['total_orders']
        customers_df['clv_score'] = customers_df['total_spent'] * np.log1p(customers_df['total_orders'])
    
    # Recency Score (días desde último login)
    if 'last_login_date' in customers_df.columns:
        current_date = datetime.now()
        customers_df['days_since_last_login'] = (current_date - customers_df['last_login_date']).dt.days
        customers_df['recency_score'] = pd.cut(customers_df['days_since_last_login'], 
                                              bins=[0, 7, 30, 90, float('inf')], 
                                              labels=[4, 3, 2, 1])
    
    # Engagement Score basado en diferentes factores
    engagement_features = []
    if 'total_orders' in customers_df.columns:
        engagement_features.append('total_orders')
    if 'days_since_last_login' in customers_df.columns:
        engagement_features.append('days_since_last_login')
    
    if engagement_features:
        # Características normalizadas para Engagement Score
        for feature in engagement_features:
            if feature in customers_df.columns:
                min_val = customers_df[feature].min()
                max_val = customers_df[feature].max()
                if max_val != min_val:
                    customers_df[f'{feature}_normalized'] = (customers_df[feature] - min_val) / (max_val - min_val)
                else:
                    customers_df[f'{feature}_normalized'] = 0
        
        # Calcular Engagement Score
        normalized_features = [f'{f}_normalized' for f in engagement_features]
        customers_df['engagement_score'] = customers_df[normalized_features].mean(axis=1)
    
    # 5. Validación de datos
    logger.info("Realizando validación de datos...")
    
    validation_errors = []
    
    # Verificar reglas de negocio
    if 'total_orders' in customers_df.columns:
        if (customers_df['total_orders'] < 0).any():
            validation_errors.append("Se encontraron cantidades de pedidos negativas")
    
    if 'total_spent' in customers_df.columns:
        if (customers_df['total_spent'] < 0).any():
            validation_errors.append("Se encontraron ingresos negativos")
    
    if 'email' in customers_df.columns:
        invalid_emails = customers_df[~customers_df['email'].str.contains('@', na=False)]
        if not invalid_emails.empty:
            validation_errors.append(f"{len(invalid_emails)} direcciones de correo inválidas encontradas")
    
    if validation_errors:
        logger.warning(f"Errores de validación: {validation_errors}")
    
    # 6. Agregar datos para Analytics
    logger.info("Agregando datos para Analytics...")
    
    # Estadísticas de segmento de clientes
    if 'customer_segment' in customers_df.columns:
        segment_stats = customers_df.groupby('customer_segment').agg({
            'customer_id': 'count',
            'total_spent': ['sum', 'mean'],
            'total_orders': ['sum', 'mean'],
            'engagement_score': 'mean'
        }).round(2)
        
        logger.info("Estadísticas de segmento de clientes creadas")
    
    # Guardar datos transformados
    transformed_json = customers_df.to_json(date_format='iso')
    task_instance.xcom_push(key='transformed_data', value=transformed_json)
    
    logger.info(f"Transformación de datos completada: {len(customers_df)} registros")
    return transformed_json

def load_data_to_warehouse(**context):
    """Carga datos transformados en el Data Warehouse"""
    logger.info("Cargando datos en Data Warehouse...")
    
    # Obtener datos transformados
    task_instance = context['task_instance']
    transformed_json = task_instance.xcom_pull(task_ids='transform_customer_data', key='transformed_data')
    
    # JSON a DataFrame
    customers_df = pd.read_json(transformed_json)
    
    # PostgreSQL Hook
    postgres_hook = PostgresHook(postgres_conn_id='postgres_default')
    
    # Crear tabla si no existe
    create_table_sql = """
    CREATE TABLE IF NOT EXISTS analytics_customer_metrics (
        customer_id VARCHAR(50) PRIMARY KEY,
        first_name VARCHAR(100),
        last_name VARCHAR(100),
        email VARCHAR(255),
        phone VARCHAR(50),
        registration_date TIMESTAMP,
        last_login_date TIMESTAMP,
        total_orders INTEGER,
        total_spent DECIMAL(10,2),
        customer_segment VARCHAR(50),
        avg_order_value DECIMAL(10,2),
        clv_score DECIMAL(10,2),
        days_since_last_login INTEGER,
        recency_score INTEGER,
        engagement_score DECIMAL(5,2),
        created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
        updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
    );
    """
    
    postgres_hook.run(create_table_sql)
    
    # Operación Upsert (Insertar o Actualizar)
    upsert_sql = """
    INSERT INTO analytics_customer_metrics (
        customer_id, first_name, last_name, email, phone, 
        registration_date, last_login_date, total_orders, total_spent,
        customer_segment, avg_order_value, clv_score, days_since_last_login,
        recency_score, engagement_score
    ) VALUES (
        %(customer_id)s, %(first_name)s, %(last_name)s, %(email)s, %(phone)s,
        %(registration_date)s, %(last_login_date)s, %(total_orders)s, %(total_spent)s,
        %(customer_segment)s, %(avg_order_value)s, %(clv_score)s, %(days_since_last_login)s,
        %(recency_score)s, %(engagement_score)s
    )
    ON CONFLICT (customer_id) 
    DO UPDATE SET
        first_name = EXCLUDED.first_name,
        last_name = EXCLUDED.last_name,
        email = EXCLUDED.email,
        phone = EXCLUDED.phone,
        registration_date = EXCLUDED.registration_date,
        last_login_date = EXCLUDED.last_login_date,
        total_orders = EXCLUDED.total_orders,
        total_spent = EXCLUDED.total_spent,
        customer_segment = EXCLUDED.customer_segment,
        avg_order_value = EXCLUDED.avg_order_value,
        clv_score = EXCLUDED.clv_score,
        days_since_last_login = EXCLUDED.days_since_last_login,
        recency_score = EXCLUDED.recency_score,
        engagement_score = EXCLUDED.engagement_score,
        updated_at = CURRENT_TIMESTAMP;
    """
    
    # Cargar datos en lotes (optimización de rendimiento)
    batch_size = 1000
    total_rows = len(customers_df)
    
    for i in range(0, total_rows, batch_size):
        batch_df = customers_df.iloc[i:i+batch_size]
        
        # Preparar datos para PostgreSQL
        batch_records = batch_df.to_dict('records')
        
        # Tratar valores NaN
        for record in batch_records:
            for key, value in record.items():
                if pd.isna(value):
                    record[key] = None
                elif isinstance(value, pd.Timestamp):
                    record[key] = value.isoformat()
        
        # Ejecutar lote
        postgres_hook.run(upsert_sql, parameters=batch_records)
        
        logger.info(f"Lote {i//batch_size + 1}/{(total_rows-1)//batch_size + 1} cargado: {len(batch_df)} filas")
    
    # Cargar también en BigQuery para Analytics
    try:
        bigquery_hook = BigQueryHook(gcp_conn_id='google_cloud_default', use_legacy_sql=False)
        
        # Configurar dataset y tabla
        project_id = Variable.get('gcp_project_id', default_var='my-project')
        dataset_id = 'analytics'
        table_id = 'customer_metrics'
        
        # Cargar DataFrame en BigQuery
        customers_df.to_gbq(
            destination_table=f'{dataset_id}.{table_id}',
            project_id=project_id,
            if_exists='replace',
            progress_bar=False
        )
        
        logger.info("Datos cargados exitosamente en BigQuery")
        
    except Exception as e:
        logger.warning(f"Carga a BigQuery falló: {e}")
    
    logger.info("Proceso de carga de datos completado")

def generate_analytics_report(**context):
    """Genera reporte de Analytics"""
    logger.info("Generando reporte de Analytics...")
    
    # PostgreSQL Hook para datos finales
    postgres_hook = PostgresHook(postgres_conn_id='postgres_default')
    
    # Ejecutar queries de Analytics
    analytics_queries = {
        'customer_segments': """
            SELECT 
                customer_segment,
                COUNT(*) as customer_count,
                AVG(total_spent) as avg_total_spent,
                AVG(total_orders) as avg_total_orders,
                AVG(clv_score) as avg_clv_score,
                AVG(engagement_score) as avg_engagement_score
            FROM analytics_customer_metrics
            GROUP BY customer_segment
            ORDER BY avg_total_spent DESC
        """,
        
        'top_customers': """
            SELECT 
                customer_id,
                first_name,
                last_name,
                total_spent,
                total_orders,
                clv_score,
                engagement_score
            FROM analytics_customer_metrics
            ORDER BY clv_score DESC
            LIMIT 10
        """,
        
        'engagement_distribution': """
            SELECT 
                CASE 
                    WHEN engagement_score >= 0.8 THEN 'High'
                    WHEN engagement_score >= 0.6 THEN 'Medium'
                    WHEN engagement_score >= 0.4 THEN 'Low'
                    ELSE 'Very Low'
                END as engagement_level,
                COUNT(*) as customer_count,
                AVG(total_spent) as avg_total_spent
            FROM analytics_customer_metrics
            GROUP BY engagement_level
            ORDER BY avg_total_spent DESC
        """,
        
        'daily_trends': """
            SELECT 
                DATE(last_login_date) as login_date,
                COUNT(*) as active_customers,
                AVG(total_spent) as avg_spend_per_customer
            FROM analytics_customer_metrics
            WHERE last_login_date >= CURRENT_DATE - INTERVAL 30 DAY
            GROUP BY DATE(last_login_date)
            ORDER BY login_date DESC
        """
    }
    
    # Recopilar resultados
    analytics_results = {}
    
    for query_name, query in analytics_queries.items():
        try:
            result_df = postgres_hook.get_pandas_df(query)
            analytics_results[query_name] = result_df
            logger.info(f"Query de Analytics '{query_name}' ejecutado: {len(result_df)} filas")
        except Exception as e:
            logger.error(f"Error en query '{query_name}': {e}")
            analytics_results[query_name] = pd.DataFrame()
    
    # Crear reporte
    report_data = {
        'generated_at': datetime.now().isoformat(),
        'total_customers': len(analytics_results.get('customer_segments', pd.DataFrame())),
        'analytics': {}
    }
    
    for query_name, result_df in analytics_results.items():
        if not result_df.empty:
            report_data['analytics'][query_name] = result_df.to_dict('records')
    
    # Guardar reporte como JSON
    report_json = json.dumps(report_data, indent=2, default=str)
    
    # Guardar reporte en S3
    try:
        s3_hook = S3Hook(aws_conn_id='aws_default')
        bucket_name = Variable.get('reports_bucket_name', default_var='analytics-reports')
        
        report_filename = f"customer_analytics_report_{datetime.now().strftime('%Y%m%d_%H%M%S')}.json"
        s3_hook.load_string(
            report_json,
            key=report_filename,
            bucket_name=bucket_name,
            replace=True
        )
        
        logger.info(f"Reporte de Analytics guardado: s3://{bucket_name}/{report_filename}")
        
    except Exception as e:
        logger.error(f"Error al guardar reporte: {e}")
    
    # Registrar métricas del reporte
    if 'customer_segments' in analytics_results:
        segments_df = analytics_results['customer_segments']
        if not segments_df.empty:
            logger.info("Distribución de segmentos de clientes:")
            for _, row in segments_df.iterrows():
                logger.info(f"  {row['customer_segment']}: {row['customer_count']} clientes, "
                           f"Ø Ingresos: €{row['avg_total_spent']:.2f}")
    
    logger.info("Generación de reporte de Analytics completada")

# Definir tareas
extract_task = PythonOperator(
    task_id='extract_customer_data',
    python_callable=extract_customer_data,
    dag=dag,
)

transform_task = PythonOperator(
    task_id='transform_customer_data',
    python_callable=transform_customer_data,
    dag=dag,
)

load_task = PythonOperator(
    task_id='load_data_to_warehouse',
    python_callable=load_data_to_warehouse,
    dag=dag,
)

report_task = PythonOperator(
    task_id='generate_analytics_report',
    python_callable=generate_analytics_report,
    dag=dag,
)

# Dependencias de tareas
extract_task >> transform_task >> load_task >> report_task

3. Visualización de datos con Matplotlib, Seaborn y Plotly

import matplotlib.pyplot as plt
import seaborn as sns
import plotly.express as px
import plotly.graph_objects as go
from plotly.subplots import make_subplots
import pandas as pd
import numpy as np
from datetime import datetime, timedelta
import warnings
warnings.filterwarnings('ignore')

# Data Visualization Demo
class DataVisualizationDemo:
    
    def __init__(self):
        """Configurar estilos y parámetros"""
        # Configuración de Matplotlib
        plt.style.use('seaborn-v0_8')
        plt.rcParams['figure.figsize'] = (12, 8)
        plt.rcParams['font.size'] = 10
        plt.rcParams['axes.titlesize'] = 14
        plt.rcParams['axes.labelsize'] = 12
        plt.rcParams['xtick.labelsize'] = 10
        plt.rcParams['ytick.labelsize'] = 10
        plt.rcParams['legend.fontsize'] = 10
        plt.rcParams['figure.titlesize'] = 16
        
        # Configuración de Seaborn
        sns.set_palette("husl")
        
        # Definir paleta de colores
        self.colors = {
            'primary': '#2E86AB',
            'secondary': '#A23B72',
            'accent': '#F18F01',
            'success': '#C73E1D',
            'warning': '#F4A261',
            'info': '#264653',
            'light': '#E9C46A',
            'dark': '#2A9D8F'
        }
        
        print("Data Visualization Demo inicializado")
    
    def create_sample_data(self):
        """Crear datos de ejemplo para visualización"""
        np.random.seed(42)
        
        # Datos de series temporales (evolución de ingresos)
        dates = pd.date_range(start='2023-01-01', end='2023-12-31', freq='D')
        base_revenue = 10000
        trend = np.linspace(0, 5000, len(dates))
        seasonal = 2000 * np.sin(2 * np.pi * np.arange(len(dates)) / 365.25)
        noise = np.random.normal(0, 500, len(dates))
        
        revenue_data = pd.DataFrame({
            'date': dates,
            'revenue': base_revenue + trend + seasonal + noise,
            'orders': np.random.poisson(100, len(dates)) + 50,
            'customers': np.random.poisson(80, len(dates)) + 30
        })
        
        # Datos de clientes
        customer_data = pd.DataFrame({
            'customer_id': range(1000),
            'age': np.random.normal(35, 10, 1000),
            'income': np.random.lognormal(10.5, 0.5, 1000),
            'total_spent': np.random.gamma(2, 100, 1000),
            'orders': np.random.poisson(5, 1000) + 1,
            'segment': np.random.choice(['Bronze', 'Silver', 'Gold', 'Platinum'], 1000, 
                                      p=[0.4, 0.3, 0.2, 0.1]),
            'registration_date': pd.date_range('2020-01-01', '2023-12-31', periods=1000)[np.random.permutation(1000)],
            'last_purchase': pd.date_range('2023-01-01', '2023-12-31', periods=1000)[np.random.permutation(1000)]
        })
        
        # Datos de productos
        products = ['Laptop', 'Smartphone', 'Tablet', 'Headphones', 'Mouse', 'Keyboard', 'Monitor', 'Webcam']
        product_data = pd.DataFrame({
            'product': np.random.choice(products, 5000),
            'category': np.random.choice(['Electronics', 'Accessories'], 5000),
            'price': np.random.uniform(10, 1000, 5000),
            'quantity_sold': np.random.poisson(10, 5000) + 1,
            'rating': np.random.uniform(3, 5, 5000),
            'month': np.random.choice(range(1, 13), 5000)
        })
        
        # Datos geográficos
        regions = ['North', 'South', 'East', 'West', 'Central']
        geo_data = pd.DataFrame({
            'region': np.random.choice(regions, 1000),
            'city': [f"{region} City {i}" for region, i in zip(
                np.random.choice(regions, 1000), 
                np.random.randint(1, 10, 1000)
            )],
            'population': np.random.lognormal(10, 1, 1000),
            'revenue': np.random.gamma(2, 50000, 1000),
            'stores': np.random.poisson(5, 1000) + 1
        })
        
        return revenue_data, customer_data, product_data, geo_data
    
    def create_time_series_visualizations(self, revenue_data):
        """Crear visualizaciones de series temporales"""
        print("Creando visualizaciones de series temporales...")
        
        fig, axes = plt.subplots(2, 2, figsize=(16, 12))
        fig.suptitle('Análisis de ingresos en series temporales', fontsize=16, fontweight='bold')
        
        # 1. Evolución diaria de ingresos
        axes[0, 0].plot(revenue_data['date'], revenue_data['revenue'], 
                       color=self.colors['primary'], linewidth=2)
        axes[0, 0].set_title('Evolución diaria de ingresos')
        axes[0, 0].set_xlabel('Fecha')
        axes[0, 0].set_ylabel('Ingresos (€)')
        axes[0, 0].grid(True, alpha=0.3)
        
        # Agregar media móvil
        revenue_data['ma_7'] = revenue_data['revenue'].rolling(window=7).mean()
        revenue_data['ma_30'] = revenue_data['revenue'].rolling(window=30).mean()
        
        axes[0, 0].plot(revenue_data['date'], revenue_data['ma_7'], 
                       color=self.colors['accent'], linewidth=1, alpha=0.7, label='Media móvil 7 días')
        axes[0, 0].plot(revenue_data['date'], revenue_data['ma_30'], 
                       color=self.colors['secondary'], linewidth=1, alpha=0.7, label='Media móvil 30 días')
        axes[0, 0].legend()
        
        # 2. Ingresos mensuales
        monthly_revenue = revenue_data.set_index('date').resample('M')['revenue'].sum()
        axes[0, 1].bar(monthly_revenue.index, monthly_revenue.values, 
                      color=self.colors['info'], alpha=0.7)
        axes[0, 1].set_title('Ingresos mensuales totales')
        axes[0, 1].set_xlabel('Mes')
        axes[0, 1].set_ylabel('Ingresos (€)')
        axes[0, 1].tick_params(axis='x', rotation=45)
        
        # 3. Ingresos vs Pedidos
        axes[1, 0].scatter(revenue_data['orders'], revenue_data['revenue'], 
                          alpha=0.6, color=self.colors['primary'])
        axes[1, 0].set_title('Ingresos vs. Pedidos')
        axes[1, 0].set_xlabel('Número de pedidos')
        axes[1, 0].set_ylabel('Ingresos (€)')
        
        # Agregar línea de tendencia
        z = np.polyfit(revenue_data['orders'], revenue_data['revenue'], 1)
        p = np.poly1d(z)
        axes[1, 0].plot(revenue_data['orders'], p(revenue_data['orders']), 
                       color=self.colors['accent'], linewidth=2)
        
        # 4. Distribución de ingresos diarios
        axes[1, 1].hist(revenue_data['revenue'], bins=30, color=self.colors['secondary'], 
                       alpha=0.7, edgecolor='black')
        axes[1, 1].set_title('Distribución de ingresos diarios')
        axes[1, 1].set_xlabel('Ingresos (€)')
        axes[1, 1].set_ylabel('Frecuencia')
        axes[1, 1].axvline(revenue_data['revenue'].mean(), color=self.colors['accent'], 
                           linestyle='--', linewidth=2, label=f'Promedio: €{revenue_data["revenue"].mean():.0f}')
        axes[1, 1].legend()
        
        plt.tight_layout()
        plt.savefig('time_series_analysis.png', dpi=300, bbox_inches='tight')
        plt.show()
        
        return fig
    
    def create_customer_analytics_visualizations(self, customer_data):
        """Crear visualizaciones de análisis de clientes"""
        print("Creando visualizaciones de análisis de clientes...")
        
        fig, axes = plt.subplots(2, 3, figsize=(18, 12))
        fig.suptitle('Panel de análisis de clientes', fontsize=16, fontweight='bold')
        
        # 1. Distribución de edades
        axes[0, 0].hist(customer_data['age'], bins=20, color=self.colors['primary'], 
                       alpha=0.7, edgecolor='black')
        axes[0, 0].set_title('Distribución de edades')
        axes[0, 0].set_xlabel('Edad')
        axes[0, 0].set_ylabel('Cantidad')
        
        # 2. Distribución de ingresos
        axes[0, 1].hist(customer_data['income'], bins=30, color=self.colors['secondary'], 
                       alpha=0.7, edgecolor='black')
        axes[0, 1].set_title('Distribución de ingresos')
        axes[0, 1].set_xlabel('Ingresos (€)')
        axes[0, 1].set_ylabel('Cantidad')
        
        # 3. Distribución de segmentos de clientes
        segment_counts = customer_data['segment'].value_counts()
        axes[0, 2].pie(segment_counts.values, labels=segment_counts.index, 
                      autopct='%1.1f%%', colors=[self.colors['primary'], self.colors['secondary'], 
                                                self.colors['accent'], self.colors['success']])
        axes[0, 2].set_title('Distribución de segmentos de clientes')
        
        # 4. Ingresos promedio por segmento
        segment_revenue = customer_data.groupby('segment')['total_spent'].mean()
        bars = axes[1, 0].bar(segment_revenue.index, segment_revenue.values, 
                              color=[self.colors['primary'], self.colors['secondary'], 
                                     self.colors['accent'], self.colors['success']])
        axes[1, 0].set_title('Ingresos promedio por segmento')
        axes[1, 0].set_xlabel('Segmento')
        axes[1, 0].set_ylabel('Ingresos promedio (€)')
        
        # Mostrar valores sobre las barras
        for bar, value in zip(bars, segment_revenue.values):
            axes[1, 0].text(bar.get_x() + bar.get_width()/2, bar.get_height() + 100,
                           f'€{value:.0f}', ha='center', va='bottom')
        
        # 5. Edad vs. Gastos totales
        scatter = axes[1, 1].scatter(customer_data['age'], customer_data['total_spent'], 
                                   alpha=0.6, c=customer_data['orders'], cmap='viridis')
        axes[1, 1].set_title('Edad vs. Gastos totales')
        axes[1, 1].set_xlabel('Edad')
        axes[1, 1].set_ylabel('Gastos totales (€)')
        plt.colorbar(scatter, ax=axes[1, 1], label='Número de pedidos')
        
        # 6. Boxplot de gastos por segmento
        segments = customer_data['segment'].unique()
        data_for_boxplot = [customer_data[customer_data['segment'] == seg]['total_spent'] 
                           for seg in segments]
        
        box_plot = axes[1, 2].boxplot(data_for_boxplot, labels=segments, patch_artist=True)
        colors = [self.colors['primary'], self.colors['secondary'], 
                 self.colors['accent'], self.colors['success']]
        
        for patch, color in zip(box_plot['boxes'], colors):
            patch.set_facecolor(color)
            patch.set_alpha(0.7)
        
        axes[1, 2].set_title('Distribución de gastos por segmento')
        axes[1, 2].set_xlabel('Segmento')
        axes[1, 2].set_ylabel('Gastos totales (€)')
        
        plt.tight_layout()
        plt.savefig('customer_analytics.png', dpi=300, bbox_inches='tight')
        plt.show()
        
        return fig
    
    def create_product_analytics_visualizations(self, product_data):
        """Crear visualizaciones de análisis de productos"""
        print("Creando visualizaciones de análisis de productos...")
        
        fig, axes = plt.subplots(2, 2, figsize=(16, 12))
        fig.suptitle('Panel de análisis de productos', fontsize=16, fontweight='bold')
        
        # 1. Productos principales por ingresos
        product_revenue = product_data.groupby('product').apply(
            lambda x: (x['price'] * x['quantity_sold']).sum()
        ).sort_values(ascending=False)
        
        bars = axes[0, 0].barh(product_revenue.index, product_revenue.values, 
                              color=self.colors['primary'])
        axes[0, 0].set_title('Productos principales por ingresos totales')
        axes[0, 0].set_xlabel('Ingresos totales (€)')
        
        # Mostrar valores sobre las barras
        for i, (bar, value) in enumerate(zip(bars, product_revenue.values)):
            axes[0, 0].text(value + max(product_revenue.values) * 0.01, i,
                           f'€{value:,.0f}', va='center')
        
        # 2. Distribución de precios por categoría
        categories = product_data['category'].unique()
        for i, category in enumerate(categories):
            category_data = product_data[product_data['category'] == category]['price']
            axes[0, 1].hist(category_data, bins=20, alpha=0.7, label=category,
                           color=self.colors['primary'] if i == 0 else self.colors['secondary'])
        
        axes[0, 1].set_title('Distribución de precios por categoría')
        axes[0, 1].set_xlabel('Precio (€)')
        axes[0, 1].set_ylabel('Frecuencia')
        axes[0, 1].legend()
        
        # 3. Calificación vs. Precio
        axes[1, 0].scatter(product_data['price'], product_data['rating'], 
                          alpha=0.6, color=self.colors['accent'])
        axes[1, 0].set_title('Calificación vs. Precio')
        axes[1, 0].set_xlabel('Precio (€)')
        axes[1, 0].set_ylabel('Calificación')
        
        # Coeficiente de correlación
        correlation = np.corrcoef(product_data['price'], product_data['rating'])[0, 1]
        axes[1, 0].text(0.05, 0.95, f'Correlación: {correlation:.3f}', 
                       transform=axes[1, 0].transAxes, va='top',
                       bbox=dict(boxstyle='round', facecolor='white', alpha=0.8))
        
        # 4. Ventas mensuales
        monthly_sales = product_data.groupby('month')['quantity_sold'].sum()
        month_names = ['Ene', 'Feb', 'Mar', 'Abr', 'May', 'Jun', 
                      'Jul', 'Ago', 'Sep', 'Oct', 'Nov', 'Dic']
        
        axes[1, 1].bar(monthly_sales.index, monthly_sales.values, 
                      color=self.colors['info'], alpha=0.7)
        axes[1, 1].set_title('Volúmenes de ventas mensuales')
        axes[1, 1].set_xlabel('Mes')
        axes[1, 1].set_ylabel('Cantidad vendida')
        axes[1, 1].set_xticks(range(1, 13))
        axes[1, 1].set_xticklabels(month_names, rotation=45)
        
        plt.tight_layout()
        plt.savefig('product_analytics.png', dpi=300, bbox_inches='tight')
        plt.show()
        
        return fig
    
    def create_geographic_visualizations(self, geo_data):
        """Crear visualizaciones geográficas"""
        print("Creando visualizaciones geográficas...")
        
        fig, axes = plt.subplots(2, 2, figsize=(16, 12))
        fig.suptitle('Panel de análisis geográfico', fontsize=16, fontweight='bold')
        
        # 1. Ingresos por región
        region_revenue = geo_data.groupby('region')['revenue'].sum().sort_values(ascending=False)
        bars = axes[0, 0].bar(region_revenue.index, region_revenue.values, 
                              color=[self.colors['primary'], self.colors['secondary'], 
                                     self.colors['accent'], self.colors['success'], self.colors['info']])
        axes[0, 0].set_title('Ingresos totales por región')
        axes[0, 0].set_ylabel('Ingresos (€)')
        
        # Mostrar valores sobre las barras
        for bar, value in zip(bars, region_revenue.values):
            axes[0, 0].text(bar.get_x() + bar.get_width()/2, bar.get_height() + max(region_revenue.values) * 0.01,
                           f'€{value:,.0f}', ha='center', va='bottom')
        
        # 2. Población vs. Ingresos
        axes[0, 1].scatter(geo_data['population'], geo_data['revenue'], 
                          alpha=0.6, s=geo_data['stores']*10, 
                          c=self.colors['primary'])
        axes[0, 1].set_title('Población vs. Ingresos')
        axes[0, 1].set_xlabel('Población')
        axes[0, 1].set_ylabel('Ingresos (€)')
        
        # 3. Densidad de tiendas por región
        region_stores = geo_data.groupby('region')['stores'].sum()
        axes[1, 0].pie(region_stores.values, labels=region_stores.index, 
                      autopct='%1.1f%%', colors=[self.colors['primary'], self.colors['secondary'], 
                                                self.colors['accent'], self.colors['success'], self.colors['info']])
        axes[1, 0].set_title('Distribución de tiendas por región')
        
        # 4. Ingresos promedio por tienda
        region_avg_revenue = geo_data.groupby('region').apply(
            lambda x: x['revenue'].sum() / x['stores'].sum()
        ).sort_values(ascending=False)
        
        bars = axes[1, 1].bar(region_avg_revenue.index, region_avg_revenue.values, 
                              color=self.colors['secondary'])
        axes[1, 1].set_title('Ingresos promedio por tienda')
        axes[1, 1].set_ylabel('Ingresos promedio por tienda (€)')
        
        # Mostrar valores sobre las barras
        for bar, value in zip(bars, region_avg_revenue.values):
            axes[1, 1].text(bar.get_x() + bar.get_width()/2, bar.get_height() + max(region_avg_revenue.values) * 0.01,
                           f'€{value:,.0f}', ha='center', va='bottom')
        
        plt.tight_layout()
        plt.savefig('geographic_analytics.png', dpi=300, bbox_inches='tight')
        plt.show()
        
        return fig
    
    def create_interactive_visualizations(self, revenue_data, customer_data, product_data):
        """Crear visualizaciones interactivas con Plotly"""
        print("Creando visualizaciones interactivas...")
        
        # 1. Series temporales interactivas
        fig1 = make_subplots(
            rows=2, cols=1,
            subplot_titles=('Ingresos diarios', 'Ingresos mensuales'),
            vertical_spacing=0.1
        )
        
        # Ingresos diarios con medias móviles
        fig1.add_trace(
            go.Scatter(
                x=revenue_data['date'],
                y=revenue_data['revenue'],
                mode='lines',
                name='Ingresos diarios',
                line=dict(color='blue', width=1)
            ),
            row=1, col=1
        )
        
        fig1.add_trace(
            go.Scatter(
                x=revenue_data['date'],
                y=revenue_data['ma_7'],
                mode='lines',
                name='Media móvil 7 días',
                line=dict(color='orange', width=2)
            ),
            row=1, col=1
        )
        
        # Ingresos mensuales
        monthly_revenue = revenue_data.set_index('date').resample('M')['revenue'].sum()
        fig1.add_trace(
            go.Bar(
                x=monthly_revenue.index,
                y=monthly_revenue.values,
                name='Ingresos mensuales',
                marker_color='lightblue'
            ),
            row=2, col=1
        )
        
        fig1.update_layout(
            title='Análisis interactivo de ingresos',
            height=600,
            showlegend=True
        )
        
        fig1.show()
        
        # 2. Análisis interactivo de clientes
        fig2 = make_subplots(
            rows=2, cols=2,
            subplot_titles=('Distribución de edades', 'Distribución de ingresos', 
                          'Ingresos por segmento', 'Edad vs. Ingresos'),
            specs=[[{"type": "histogram"}, {"type": "histogram"}],
                   [{"type": "bar"}, {"type": "scatter"}]]
        )
        
        # Distribución de edades
        fig2.add_trace(
            go.Histogram(
                x=customer_data['age'],
                nbinsx=20,
                name='Edad',
                marker_color='lightgreen'
            ),
            row=1, col=1
        )
        
        # Distribución de ingresos
        fig2.add_trace(
            go.Histogram(
                x=customer_data['income'],
                nbinsx=30,
                name='Ingresos',
                marker_color='lightcoral'
            ),
            row=1, col=2
        )
        
        # Ingresos por segmento
        segment_revenue = customer_data.groupby('segment')['total_spent'].mean()
        fig2.add_trace(
            go.Bar(
                x=segment_revenue.index,
                y=segment_revenue.values,
                name='Ingresos promedio',
                marker_color='lightblue'
            ),
            row=2, col=1
        )
        
        # Edad vs. Ingresos
        fig2.add_trace(
            go.Scatter(
                x=customer_data['age'],
                y=customer_data['total_spent'],
                mode='markers',
                name='Clientes',
                marker=dict(
                    size=8,
                    color=customer_data['orders'],
                    colorscale='Viridis',
                    showscale=True,
                    colorbar=dict(title="Pedidos")
                )
            ),
            row=2, col=2
        )
        
        fig2.update_layout(
            title='Análisis interactivo de clientes',
            height=600,
            showlegend=False
        )
        
        fig2.show()
        
        # 3. Visualización 3D interactiva
        fig3 = go.Figure(data=[go.Scatter3d(
            x=customer_data['age'],
            y=customer_data['income'],
            z=customer_data['total_spent'],
            mode='markers',
            marker=dict(
                size=5,
                color=customer_data['orders'],
                colorscale='Viridis',
                showscale=True,
                colorbar=dict(title="Pedidos")
            ),
            text=[f"Cliente {i}<br>Edad: {age}<br>Ingresos: {inc:,.0f}€<br>Gastos: {spent:,.0f}€" 
                   for i, (age, inc, spent) in enumerate(zip(customer_data['age'], 
                                                          customer_data['income'], 
                                                          customer_data['total_spent']))],
            hovertemplate='%{text}<extra></extra>'
        )])
        
        fig3.update_layout(
            title='Análisis 3D de clientes (Edad, Ingresos, Gastos)',
            scene=dict(
                xaxis_title='Edad',
                yaxis_title='Ingresos (€)',
                zaxis_title='Gastos (€)'
            ),
            height=600
        )
        
        fig3.show()
        
        return fig1, fig2, fig3
    
    def create_dashboard_summary(self, revenue_data, customer_data, product_data, geo_data):
        """Crear panel de control resumido"""
        print("Creando panel de control resumido...")
        
        # Calcular métricas clave
        total_revenue = revenue_data['revenue'].sum()
        avg_daily_revenue = revenue_data['revenue'].mean()
        total_customers = len(customer_data)
        avg_customer_value = customer_data['total_spent'].mean()
        total_orders = customer_data['orders'].sum()
        top_product = product_data.groupby('product').apply(
            lambda x: (x['price'] * x['quantity_sold']).sum()
        ).idxmax()
        
        # Crear panel
        fig = plt.figure(figsize=(20, 16))
        
        # Título
        fig.suptitle('Panel de análisis en ciencia de datos', fontsize=20, fontweight='bold')
        
        # Métricas clave
        metrics_text = f"""
        MÉTRICAS CLAVE
        ─────────────────────
        Ingresos totales: €{total_revenue:,.0f}
        Ingresos promedio diario: €{avg_daily_revenue:,.0f}
        Total de clientes: {total_customers:,}
        Valor promedio de cliente: €{avg_customer_value:,.0f}
        Total de pedidos: {total_orders:,}
        Producto principal: {top_product}
        """
        
        plt.figtext(0.02, 0.95, metrics_text, fontsize=12, fontfamily='monospace',
                   verticalalignment='top', bbox=dict(boxstyle='round', facecolor='lightgray', alpha=0.8))
        
        # Subpaneles para diferentes análisis
        gs = fig.add_gridspec(3, 3, hspace=0.3, wspace=0.3, 
                           left=0.15, right=0.95, top=0.85, bottom=0.05)
        
        # 1. Tendencia de ingresos (arriba izquierda)
        ax1 = fig.add_subplot(gs[0, 0])
        monthly_revenue = revenue_data.set_index('date').resample('M')['revenue'].sum()
        ax1.plot(monthly_revenue.index, monthly_revenue.values, 
                color=self.colors['primary'], linewidth=2)
        ax1.set_title('Tendencia de ingresos mensuales')
        ax1.tick_params(axis='x', rotation=45)
        
        # 2. Segmentos de clientes (arriba centro)
        ax2 = fig.add_subplot(gs[0, 1])
        segment_counts = customer_data['segment'].value_counts()
        ax2.pie(segment_counts.values, labels=segment_counts.index, autopct='%1.1f%%',
               colors=[self.colors['primary'], self.colors['secondary'], 
                      self.colors['accent'], self.colors['success']])
        ax2.set_title('Segmentos de clientes')
        
        # 3. Productos principales (arriba derecha)
        ax3 = fig.add_subplot(gs[0, 2])
        product_revenue = product_data.groupby('product').apply(
            lambda x: (x['price'] * x['quantity_sold']).sum()
        ).sort_values(ascending=False).head(5)
        
        bars = ax3.barh(range(len(product_revenue)), product_revenue.values, 
                       color=self.colors['info'])
        ax3.set_yticks(range(len(product_revenue)))
        ax3.set_yticklabels(product_revenue.index)
        ax3.set_title('5 productos principales')
        ax3.set_xlabel('Ingresos (€)')
        
        # 4. Distribución de edades (medio izquierda)
        ax4 = fig.add_subplot(gs[1, 0])
## Flujo de Trabajo en Data Science

### Proceso Típico de Data Science
```mermaid
graph TD
    A[Problem Definition] --> B[Data Collection]
    B --> C[Data Cleaning]
    C --> D[Exploratory Analysis]
    D --> E[Feature Engineering]
    E --> F[Model Building]
    F --> G[Model Evaluation]
    G --> H[Deployment]
    H --> I[Monitoring]
    I --> A

Metodología CRISP-DM

  1. Business Understanding: Comprender los objetivos del negocio
  2. Data Understanding: Explorar los datos
  3. Data Preparation: Limpiar y preparar los datos
  4. Modeling: Desarrollar modelos
  5. Evaluation: Evaluar los modelos
  6. Deployment: Desplegar los modelos
  7. Monitoring: Supervisar el rendimiento

Tecnologías Big Data

Ecosistema Hadoop

ComponenteFunciónAlternativas
HDFSAlmacenamiento distribuidoAmazon S3, Google Cloud Storage
MapReduceProcesamiento por lotesApache Spark, Apache Flink
HiveSQL en HadoopPresto, Apache Impala
HBaseBase de datos NoSQLCassandra, MongoDB
ZookeeperCoordinaciónConsul, etcd

Spark vs. MapReduce

CaracterísticaSparkMapReduce
VelocidadEn memoriaBasado en disco
Facilidad de usoAPIs de alto nivelBajo nivel
StreamingNativoSolo por lotes
Soporte MLMLlibExterno
Uso de recursosMayorMenor

Componentes del Pipeline de Machine Learning

Etapas del Pipeline de Datos

# Ejemplo de Pipeline ML con Scikit-learn
from sklearn.pipeline import Pipeline
from sklearn.preprocessing import StandardScaler, OneHotEncoder
from sklearn.compose import ColumnTransformer
from sklearn.ensemble import RandomForestRegressor
from sklearn.model_selection import train_test_split

# Preprocesamiento
numeric_features = ['age', 'income', 'orders']
categorical_features = ['segment', 'region']

numeric_transformer = Pipeline(steps=[
    ('scaler', StandardScaler())
])

categorical_transformer = Pipeline(steps=[
    ('onehot', OneHotEncoder(handle_unknown='ignore'))
])

preprocessor = ColumnTransformer(
    transformers=[
        ('num', numeric_transformer, numeric_features),
        ('cat', categorical_transformer, categorical_features)
    ])

# Pipeline completo
ml_pipeline = Pipeline(steps=[
    ('preprocessor', preprocessor),
    ('regressor', RandomForestRegressor(n_estimators=100))
])

Prácticas MLOps

  • Control de versiones: Git para código, DVC para datos
  • Tracking de experimentos: MLflow, Weights & Biases
  • Registro de modelos: MLflow Model Registry
  • CI/CD: GitHub Actions, Jenkins
  • Monitoreo: Prometheus, Grafana
  • Detección de drift: Evidently AI, NannyML

Principios de Visualización de Datos

Guía de Selección de Gráficos

Tipo de datosVisualizaciónPropósito
Series de tiempoGráfico de líneasTendencias a lo largo del tiempo
CategóricoGráfico de barrasComparaciones
DistribuciónHistogramaFrecuencias
RelaciónGráfico de dispersiónCorrelaciones
ComposiciónGráfico circularProporciones
GeográficoMapaUbicaciones

Teoría del Color para Visualización de Datos

  • Colores primarios: Azul, rojo, amarillo
  • Colores secundarios: Verde, naranja, púrpura
  • Psicología del color: Azul=confianza, rojo=peligro, verde=éxito
  • Contrastes: Claro/oscuro para legibilidad
  • Daltonismo: 8-10% de la población

Mejores Prácticas ETL

Validación de Calidad de Datos

def validate_data_quality(df):
    """Validar la calidad de los datos"""
    quality_report = {
        'total_rows': len(df),
        'missing_values': df.isnull().sum().to_dict(),
        'duplicate_rows': df.duplicated().sum(),
        'data_types': df.dtypes.to_dict(),
        'outliers': detect_outliers(df)
    }
    return quality_report

def detect_outliers(df, column):
    """Detectar valores atípicos usando el método IQR"""
    Q1 = df[column].quantile(0.25)
    Q3 = df[column].quantile(0.75)
    IQR = Q3 - Q1
    lower_bound = Q1 - 1.5 * IQR
    upper_bound = Q3 + 1.5 * IQR
    outliers = df[(df[column] < lower_bound) | (df[column] > upper_bound)]
    return len(outliers)

Optimización de Rendimiento

  • Procesamiento paralelo: Multiprocessing, Dask
  • Gestión de memoria: Chunking, Generators
  • Indexación: Índices de bases de datos
  • Caché: Redis, Memcached
  • Procesamiento por lotes: Operaciones masivas

Técnicas de Análisis

Métodos Estadísticos

  • Estadística descriptiva: Media, mediana, moda, desviación estándar
  • Estadística inferencial: Pruebas de hipótesis, intervalos de confianza
  • Análisis de correlación: Pearson, Spearman, Kendall
  • Análisis de regresión: Lineal, logística, polinómica
  • Series de tiempo: ARIMA, Prophet, LSTM

Análisis Avanzados

  • Clustering: K-Means, jerárquico, DBSCAN
  • Clasificación: Árboles de decisión, Random Forest, SVM
  • Detección de anomalías: Isolation Forest, One-Class SVM
  • Reducción de dimensionalidad: PCA, t-SNE, UMAP
  • Reglas de asociación: Apriori, FP-Growth

Ventajas e Inconvenientes

Ventajas de Data Science

  • Decisiones basadas en datos: Mejores decisiones empresariales
  • Reconocimiento de patrones: Descubrir patrones ocultos
  • Análisis predictivo: Predecir tendencias futuras
  • Optimización de procesos: Mejorar la eficiencia
  • Ventaja competitiva: Obtener ventajas frente a la competencia

Inconvenientes

  • Calidad de datos: Depende de la calidad de los datos
  • Complejidad: Algoritmos y herramientas complejas
  • Preocupaciones de privacidad: Protección de datos y ética
  • Consumo intensivo de recursos: Potencia computacional y almacenamiento
  • Interpretabilidad: Modelos de caja negra

Preguntas Frecuentes de Examen

  1. ¿Cuál es la diferencia entre Big Data y datos tradicionales? Big Data se caracteriza por las 3V (volumen, velocidad, variedad) y requiere tecnologías especiales para su procesamiento.

  2. ¡Explica el proceso ETL! ETL significa Extract (extraer datos), Transform (limpiar y transformar datos) y Load (cargar datos en el sistema destino).

  3. ¿Cuándo se utiliza cada tipo de gráfico? Series de tiempo para tendencias, gráficos de barras para comparaciones, histogramas para distribuciones, gráficos de dispersión para relaciones.

  4. ¿Cuál es el propósito de los pipelines de Machine Learning? Los pipelines de ML automatizan y estandarizan todo el flujo de trabajo, desde la preparación de datos hasta el despliegue del modelo.

Fuentes Principales

  1. https://spark.apache.org/
  2. https://airflow.apache.org/
  3. https://matplotlib.org/
  4. https://plotly.com/python/
Volver al blog
Share:

Entradas relacionadas