Spark
03 / 03

Streaming, MLlib & Performance Tuning

Apache Spark: Streaming, MLlib & Performance Tuning

Structured Streaming

# Structured Streaming: same DataFrame API for real-time data

# Read from Kafka
stream_df = spark.readStream     .format("kafka")     .option("kafka.bootstrap.servers", "broker:9092")     .option("subscribe", "orders")     .option("startingOffsets", "latest")     .load()

# Parse JSON payload
from pyspark.sql.types import StructType, StringType, DoubleType
schema = StructType().add("order_id", StringType()).add("amount", DoubleType())

parsed = stream_df.select(
    F.from_json(F.col("value").cast("string"), schema).alias("data")
).select("data.*")

# Windowed aggregation (5-minute tumbling window)
windowed = parsed     .withWatermark("event_time", "10 minutes")     .groupBy(F.window("event_time", "5 minutes"))     .agg(F.sum("amount").alias("total_revenue"))

# Write to sink
query = windowed.writeStream     .outputMode("update")     .format("console")     .trigger(processingTime="30 seconds")     .start()

# Write to Kafka
query = parsed.select(
    F.col("order_id").cast("string").alias("key"),
    F.to_json(F.struct("*")).alias("value")
).writeStream     .format("kafka")     .option("kafka.bootstrap.servers", "broker:9092")     .option("topic", "processed-orders")     .option("checkpointLocation", "s3://bucket/checkpoints/")     .start()

query.awaitTermination()

MLlib

from pyspark.ml.feature import VectorAssembler, StandardScaler, StringIndexer
from pyspark.ml.classification import LogisticRegression, RandomForestClassifier
from pyspark.ml.regression import LinearRegression, GBTRegressor
from pyspark.ml import Pipeline
from pyspark.ml.evaluation import BinaryClassificationEvaluator
from pyspark.ml.tuning import CrossValidator, ParamGridBuilder

# Feature engineering
indexer = StringIndexer(inputCol="category", outputCol="category_idx")
assembler = VectorAssembler(
    inputCols=["amount", "category_idx", "user_age"],
    outputCol="features"
)
scaler = StandardScaler(inputCol="features", outputCol="scaled_features")

# Model
lr = LogisticRegression(featuresCol="scaled_features", labelCol="label",
    maxIter=100, regParam=0.01)

# Pipeline
pipeline = Pipeline(stages=[indexer, assembler, scaler, lr])

# Train/test split
train, test = df.randomSplit([0.8, 0.2], seed=42)
model = pipeline.fit(train)
predictions = model.transform(test)

# Evaluate
evaluator = BinaryClassificationEvaluator(labelCol="label")
auc = evaluator.evaluate(predictions)
print(f"AUC: {auc:.4f}")

# Cross-validation + hyperparameter tuning
paramGrid = ParamGridBuilder()     .addGrid(lr.regParam, [0.01, 0.1, 1.0])     .addGrid(lr.maxIter, [50, 100])     .build()

cv = CrossValidator(estimator=pipeline, estimatorParamMaps=paramGrid,
    evaluator=evaluator, numFolds=5)
cv_model = cv.fit(train)

# Save/load
model.save("s3://bucket/models/churn-v1")
from pyspark.ml import PipelineModel
loaded = PipelineModel.load("s3://bucket/models/churn-v1")

Performance Tuning

  • Use Parquet or ORC — columnar formats enable predicate pushdown and column pruning.

  • Adaptive Query Execution (AQE): enable with spark.sql.adaptive.enabled=true — auto-optimizes joins and partitions at runtime.

  • Broadcast joins: use F.broadcast(small_df) for tables < 10MB; avoids shuffle. Threshold: spark.sql.autoBroadcastJoinThreshold.

  • Avoid UDFs: Python UDFs serialize/deserialize every row. Use built-in functions (F.*) or Pandas UDFs for vectorized execution.

  • Data skew: add random salt to skewed keys for aggregations; use salted join technique.

  • Partition pruning: always filter on partition columns (year, month) early; saves reading entire dataset.

  • spark.sql.shuffle.partitions: default 200 is too high for small jobs, too low for large. Rule: 2-3x number of cores.

  • Cache strategically: only cache DataFrames reused multiple times. Cache after filtering/aggregating, not before.

Keep your own version of these notes — editable, searchable, and organised by your stack.

Start free