from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField, StringType, DoubleType
import datetime
# ── 1. SESION OPTIMIZADA ─────────────────────────────────────────────────────
spark = (SparkSession.builder
.appName("daily_sales_aggregation_v2")
# AQE: adapta particiones, detecta y divide particiones skewed
.config("spark.sql.adaptive.enabled", "true")
.config("spark.sql.adaptive.coalescePartitions.enabled", "true")
.config("spark.sql.adaptive.skewJoin.enabled", "true")
.config("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")
# Broadcast: product_catalog 450MB < 512MB threshold
.config("spark.sql.autoBroadcastJoinThreshold", "512m")
# Memoria ejecutores r5.4xlarge (128GB, 16 cores → 3 ejecutores)
.config("spark.executor.memory", "24g")
.config("spark.executor.memoryOverhead", "4g")
.config("spark.executor.cores", "5")
.config("spark.memory.fraction", "0.6")
.config("spark.memory.storageFraction", "0.4")
# Serializacion Kryo: 30-40% menos memoria en shuffle
.config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.config("spark.sql.execution.arrow.pyspark.enabled", "true")
# Particiones optimas para 47 GB de datos (128 MB/partition ~ 376)
.config("spark.sql.shuffle.partitions", "auto")
.config("spark.sql.files.maxPartitionBytes", "128m")
# Compresion LZ4 para shuffle (mas rapida que snappy)
.config("spark.io.compression.codec", "lz4")
.config("spark.shuffle.compress", "true")
.getOrCreate())
processing_date = datetime.date.today().strftime("%Y-%m-%d")
S3_BASE = "s3://dataflow-retail-datalake"
# ── 2. LECTURA EFICIENTE (predicate pushdown + column pruning) ────────────────
orders = (spark.read
.format("parquet")
.option("mergeSchema", "false")
.load(f"{S3_BASE}/raw/orders/ingestion_date={processing_date}/")
.select("order_id", "store_id", "customer_id",
"order_ts", "channel"))
order_items = (spark.read
.format("parquet")
.option("mergeSchema", "false")
.load(f"{S3_BASE}/raw/order_items/ingestion_date={processing_date}/")
.select("order_id", "product_id", "quantity", "unit_price", "unit_cost"))
# Tablas pequenas: broadcast sin shuffle
product_catalog = (spark.read
.parquet(f"{S3_BASE}/reference/product_catalog/")
.select("product_id", "category", "subcategory", "brand"))
store_meta = (spark.read
.parquet(f"{S3_BASE}/reference/store_metadata/")
.select("store_id", "region", "city", "store_type"))
# ── 3. TRANSFORMACIONES (sin UDFs Python) ────────────────────────────────────
# Enriquecer order_items con catalog (broadcast — sin shuffle)
items_enriched = order_items.join(
F.broadcast(product_catalog), "product_id", "left"
).withColumns({
"revenue": F.col("quantity") * F.col("unit_price"),
"cost_total": F.col("quantity") * F.col("unit_cost"),
# Margen calculado con Catalyst (antes era UDF Python)
"margin_pct": F.when(
(F.col("unit_price") > 0),
(F.col("unit_price") - F.col("unit_cost")) / F.col("unit_price")
).otherwise(F.lit(0.0))
})
# Unir con orders y metadata de tienda
sales_full = (items_enriched
.join(orders, "order_id", "inner")
.join(F.broadcast(store_meta), "store_id", "left"))
# ── 4. AGREGACIONES EFICIENTES ───────────────────────────────────────────────
daily_summary = (sales_full
.groupBy("region", "channel", "category", "brand")
.agg(
F.count("order_id").alias("num_orders"),
F.sum("quantity").alias("units_sold"),
F.sum("revenue").alias("total_revenue"),
F.sum("cost_total").alias("total_cost"),
F.avg("margin_pct").alias("avg_margin"),
F.approx_count_distinct("customer_id").alias("unique_customers")
)
.withColumn("processing_date", F.lit(processing_date))
.withColumn("gross_margin",
F.col("total_revenue") - F.col("total_cost")))
# ── 5. ESCRITURA PARTICIONADA ─────────────────────────────────────────────────
(daily_summary.write
.format("parquet")
.option("compression", "snappy")
.partitionBy("processing_date", "region")
.mode("overwrite")
.save(f"{S3_BASE}/curated/daily_sales_summary/"))
spark.stop()