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
- Data Collection: Fuentes y recopilación de datos
- Data Storage: Almacenamiento y organización
- Data Processing: Limpieza y transformación
- Data Analysis: Análisis estadístico y exploratorio
- Machine Learning: Modelos y algoritmos
- Data Visualization: Representación visual
- Model Deployment: Despliegue en producción
- 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
- Business Understanding: Comprender los objetivos del negocio
- Data Understanding: Explorar los datos
- Data Preparation: Limpiar y preparar los datos
- Modeling: Desarrollar modelos
- Evaluation: Evaluar los modelos
- Deployment: Desplegar los modelos
- Monitoring: Supervisar el rendimiento
Tecnologías Big Data
Ecosistema Hadoop
| Componente | Función | Alternativas |
|---|---|---|
| HDFS | Almacenamiento distribuido | Amazon S3, Google Cloud Storage |
| MapReduce | Procesamiento por lotes | Apache Spark, Apache Flink |
| Hive | SQL en Hadoop | Presto, Apache Impala |
| HBase | Base de datos NoSQL | Cassandra, MongoDB |
| Zookeeper | Coordinación | Consul, etcd |
Spark vs. MapReduce
| Característica | Spark | MapReduce |
|---|---|---|
| Velocidad | En memoria | Basado en disco |
| Facilidad de uso | APIs de alto nivel | Bajo nivel |
| Streaming | Nativo | Solo por lotes |
| Soporte ML | MLlib | Externo |
| Uso de recursos | Mayor | Menor |
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 datos | Visualización | Propósito |
|---|---|---|
| Series de tiempo | Gráfico de líneas | Tendencias a lo largo del tiempo |
| Categórico | Gráfico de barras | Comparaciones |
| Distribución | Histograma | Frecuencias |
| Relación | Gráfico de dispersión | Correlaciones |
| Composición | Gráfico circular | Proporciones |
| Geográfico | Mapa | Ubicaciones |
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
-
¿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.
-
¡Explica el proceso ETL! ETL significa Extract (extraer datos), Transform (limpiar y transformar datos) y Load (cargar datos en el sistema destino).
-
¿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.
-
¿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
- https://spark.apache.org/
- https://airflow.apache.org/
- https://matplotlib.org/
- https://plotly.com/python/


