Big data analytics platform with Spark and Kafka integration
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
val spark = SparkSession.builder
.appName("SparkAnalytics")
.getOrCreate()
import spark.implicits._
val events = spark.read
.format("kafka")
.option("kafka.bootstrap.servers", "localhost:9092")
.option("subscribe", "events")
.load()
val metrics = events
.select(from_json(col("value").cast("string"), schema).alias("data"))
.select("data.*")
.groupBy(window($"timestamp", "1 hour"), $"category")
.agg(count("*").alias("event_count"))
.orderBy($"window.start")
metrics.show()