"""
ETL distribuido con Dask — EcoMarket Analytics SL
Pipeline: 5.400 CSV diarios → Parquet analítico
Autor: CULTIVA IA | dask>=2026.3.0 | Python 3.12
"""
import dask
import dask.dataframe as dd
from dask.distributed import Client
from pathlib import Path
import logging
logging.basicConfig(level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s")
log = logging.getLogger("ecomarket-etl")
# ────────────────────────────────────────────────────────────
# CONFIGURACIÓN
# ────────────────────────────────────────────────────────────
RAW_GLOB = "data/ventas/ventas_*.csv"
OUT_DIR = Path("output/analisis")
CHUNK_MB = 100 # chunk objetivo (~100 MB/partición)
N_WORKERS = 16 # cores disponibles
MEM_LIMIT = "2GB" # por worker
DTYPE_MAP = {
"tienda_id": "category",
"producto_sku": "category",
"categoria": "category",
"canal": "category",
"precio_unitario": "float32",
"cantidad": "int32",
"descuento_pct": "float32",
}
# ────────────────────────────────────────────────────────────
# CLIENTE DISTRIBUIDO (dashboard: http://localhost:8787)
# ────────────────────────────────────────────────────────────
def make_client() -> Client:
client = Client(
n_workers=N_WORKERS,
threads_per_worker=1,
memory_limit=MEM_LIMIT,
)
log.info(f"Dashboard → {client.dashboard_link}")
return client
# ────────────────────────────────────────────────────────────
# EXTRACT — leer todos los CSVs en paralelo
# ────────────────────────────────────────────────────────────
def extract(glob: str) -> dd.DataFrame:
ddf = dd.read_csv(
glob,
dtype=DTYPE_MAP,
parse_dates=["fecha"],
blocksize=f"{CHUNK_MB}MiB", # ≈100 MB por partición
assume_missing=True,
)
log.info(f"Particiones detectadas: {ddf.npartitions}")
return ddf
# ────────────────────────────────────────────────────────────
# TRANSFORM — limpieza + feature engineering (lazy)
# ────────────────────────────────────────────────────────────
def transform(ddf: dd.DataFrame) -> dd.DataFrame:
ddf = ddf.dropna(subset=["fecha", "tienda_id", "precio_unitario"])
# Columnas derivadas
ddf["revenue"] = (
ddf["precio_unitario"] * ddf["cantidad"]
* (1 - ddf["descuento_pct"] / 100)
)
ddf["semana"] = ddf["fecha"].dt.isocalendar().week.astype("int32")
ddf["año"] = ddf["fecha"].dt.year.astype("int16")
ddf["mes"] = ddf["fecha"].dt.month.astype("int8")
return ddf
# ────────────────────────────────────────────────────────────
# LOAD — 4 tablas analíticas a Parquet
# ────────────────────────────────────────────────────────────
def load(ddf: dd.DataFrame) -> None:
OUT_DIR.mkdir(parents=True, exist_ok=True)
# A) Revenue por categoría × mes
agg_cat = (
ddf.groupby(["año", "mes", "categoria"])
["revenue", "cantidad"]
.agg(["sum", "count"])
)
# B) Revenue por tienda × semana
agg_store = (
ddf.groupby(["año", "semana", "tienda_id"])
["revenue"].agg(["sum", "mean", "count"])
)
# C) Top SKUs (sin cortar en Python: map_partitions)
top_sku = (
ddf.groupby("producto_sku")["revenue"].sum()
.nlargest(1000)
)
# D) Canal online vs offline × mes
canal_pivot = (
ddf.groupby(["año", "mes", "canal"])
["revenue"].sum()
)
# Compute todas las agregaciones en un solo pass
(r_cat, r_store, r_sku, r_canal) = dask.compute(
agg_cat, agg_store, top_sku, canal_pivot
)
# Guardar resultados
r_cat.to_parquet(OUT_DIR / "revenue_categoria.parquet")
r_store.to_parquet(OUT_DIR / "revenue_tienda_semana.parquet")
r_sku.to_frame().to_parquet(OUT_DIR / "top_skus.parquet")
r_canal.to_frame().to_parquet(OUT_DIR / "canal_mix.parquet")
log.info("✓ 4 tablas Parquet escritas en output/analisis/")
# ────────────────────────────────────────────────────────────
# MAIN
# ────────────────────────────────────────────────────────────
if __name__ == "__main__":
client = make_client()
try:
ddf = extract(RAW_GLOB)
ddf = transform(ddf)
load(ddf)
finally:
client.close()