I have application which runs query on Kafka topics with the schema specified, Below is my code :
SparkSession spark = SparkSession.builder()
.appName("Spark-Kafka-Integration")
.config("spark.master", "local")
.getOrCreate();
Dataset<Row> df = spark
.readStream()
.format("kafka")
.option("kafka.bootstrap.servers", "abc:9092,bcs:9092")
.option("subscribe","topic")
.option("auto.offset.reset", "latest")
.option("checkpointLocation", "/tmp")
.load();
// Mapping it to the schema
Dataset<Row> ds2 = df.select( from_json(col("value").cast("string") , Kafkaschema).as("rows"),col("timestamp"));
ds2.createOrReplaceTempView("ds2");
// Making a Row having timestamp and the values
Dataset<Row> ds3 = spark.sql("select rows.* , timestamp from ds2 ");
ds3.createOrReplaceTempView("table");
Dataset<Row> result2 = spark.sql(query.getQuery());
This runs fine, now I have view table which will have all columns and timestamp. Then I can run SQL like Select column1 , column2 from table group by window(timestamp,'1 minutes'),column1 , column2
My Question :
Is this is an efficient way to do it ? Because if I have multiple topics i.e .option("subscribe","topic1,topics2,...") then I have to create multiple data frame in order to run Join Query on them and how I can handle timestamp column ?
In case of multiple topics I will have the following code :
Dataset<Row> df = spark
.readStream()
.format("kafka")
.option("kafka.bootstrap.servers", "abc:9092,bcs:9092")
.option("subscribe","topic1, topic2,....topicn")
.option("auto.offset.reset", "latest")
.option("checkpointLocation", "/tmp")
.load();
Dataset<Row> ds = df.select( from_json(col("value").cast("string") , Kafkaschema).as("rows"),col("timestamp")).where("topic=topic1");
Dataset<Row> ds1 = df.select( from_json(col("value").cast("string") , Kafkaschema).as("rows"),col("timestamp")).where("topic=topic2");
.... so on and have to same for other data frame