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