Mapping Kafka to Spark dataFrame with the Schema

Viewed 170

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
0 Answers
Related