36 items
Coalesce Won't Increase Spark Partitions
Code QuizDistributed Feature Scaling and Shuffles in Spark
QuizImputation Leakage and Downstream Bias
Code QuizImputation Leakage and Downstream Bias
QuizShuffles and Joins in Distributed Feature Engineering
FlashcardImputation Bias and Missingness Mechanisms
FlashcardPartitioning, Shuffles, and Joins in Spark
Slides / VideoImputation Bias in Missing Data
Slides / VideoImputing Sentinel-Coded Missing Values
Code QuizMissing Data Mechanisms and Imputation
QuizScaling Before the Train/Test Split
Code QuizPreventing Leakage in Scaled Feature Pipelines
QuizFeature Pipelines at Scale
FlashcardMissing Data, Outliers & Data Quality
FlashcardFeature Pipelines at Scale
Slides / VideoMissing Data, Outliers & Data Quality
Slides / VideoMean Imputation With Sentinel Values
Code QuizPreventing Data Leakage in Imputation
QuizScaling Train and Test Data
Code QuizEncoding High-Cardinality Categorical Features
QuizFeature Engineering & Encoding Essentials
FlashcardFeature Engineering & Encoding Essentials
Slides / VideoDropping Missing Rows in Pandas
Code QuizHandling Missing Data in Pandas
QuizFiltering a DataFrame with two conditions
Code QuizFiltering Rows in Pandas
QuizHandling Missing Data & Duplicates
Slides / VideoPandas: Load, Filter & Select Data
Slides / VideoHandling Missing Data & Duplicates
FlashcardPandas: Loading, Filtering & Selecting Data
FlashcardEncoding, Scaling & Transformations
Slides / VideoHandling Missing Data & Outliers
Slides / VideoEncoding, Scaling & Transformations
FlashcardHandling Missing Data & Outliers
FlashcardOne-Hot Encoding a Category Column
Code QuizDetecting Outliers with IQR
Code QuizCoalesce Won't Increase Spark Partitions
A subtle Spark bug where coalesce fails to increase parallelism before a scaling join.
from pyspark.sql import functions as F
from pyspark.sql.functions import broadcast
# Source table read from a few large files -> ~8 input partitions
df = spark.read.parquet("s3://events/raw")
# Heavy filter drops ~90% of rows but keeps skewed key distribution
df = df.filter(F.col("active") == True)
# We want MORE parallelism for the feature-scaling stage below,
# so bump partitions up to 400 before the join.
df = df.coalesce(400)
# Per-key scaling stats to standardize the 'amount' feature
stats = (df.groupBy("user_id")
.agg(F.mean("amount").alias("mu"),
F.stddev("amount").alias("sd")))
scaled = (df.join(broadcast(stats), on="user_id")
.withColumn("amount_z", (F.col("amount") - F.col("mu")) / F.col("sd")))
scaled.write.parquet("s3://events/scaled")The scaling stage stays slow and single-threaded despite requesting 400 partitions. What is the bug?