End-to-end data engineering platform with Kafka, Spark, and PostgreSQL
from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col
spark = SparkSession.builder \
.appName("DataEngineering") \
.config("spark.jars.packages",
"org.apache.spark:spark-sql-kafka") \
.getOrCreate()
df = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "raw_events") \
.load()
transformed = df.select(
from_json(col("value").cast("string"), schema).alias("data")
).select("data.*")
transformed.writeStream \
.format("jdbc") \
.option("url", "jdbc:postgresql://localhost/analytics") \
.start()