how to determine where the Kafka consumption is made on Spark

Viewed 136

We tried two MlLib transformers for Kafka consuming: one using Structured Batch Query Stream like here, and one using normal Kafka consumers. The normal consumers read into a list that gets converted to a dataframe

We created two MlLib pipelines starting with an empty dataframe. The first transformer of these pipes did the reading from Kafka. We had a pipe for each kafka transformer type.

then ran the pipes: 1st with normal kafka consumers, 2nd with Spark consumer, 3rd with normal consumers again.

The spark config had 4 executors, with 1 core each: spark.executor.cores": "1", "spark.executor.instances": 4,

Questions are:

A. where is the consumption made? On the executors or on the driver? according to driver UI, it looks like in both cases the executors did all the work - the driver did not pass any data and 4 executors got created.

B. Why do we have a different number of executors running? In the 1st run, with normal consumers, we see 4 executors working; In the 2nd run, with spark Kafka connector, 1 executor; In 3rd run with normal consumers, 1 executor but 2 cores?

you`ll see the driver's UI attached at the bottom.

this is the relevant code: Normal Consumer:

var kafkaConsumer: KafkaConsumer[String, String] = null
val readMessages = () => {
    for (record <- records) {
        recordList.append(record.value())
    }
}

kafkaConsumer.subscribe(util.Arrays.asList($(topic)))
readMessages()

var df = recordList.toDF
kafkaConsumer.close()
val json_schema =     
df.sparkSession.read.json(df.select("value").as[String]).schema
df = df.select(from_json(col("value"), json_schema).as("json"))
df = df.select(col("json.*"))

Spark consumer:

val records = dataset
  .sparkSession
  .read
  .format("kafka")
  .option("kafka.bootstrap.servers", $(url))
  .option("subscribe", $(this.topic))
  .option("kafkaConsumer.pollTimeoutMs", s"${$(timeoutMs)}")
  .option("startingOffsets", $(startingOffsets))
  .option("endingOffsets", $(endingOffsets))
  .load
  .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
  .as[(String, String)]

OmnixLogger.warn(uid, "Executor ID AFTER polling: " + SparkEnv.get.executorId)

val json_schema = records.sparkSession.read.json(records.select("value").as[String]).schema
var df: DataFrame = records.select(from_json(col("value"), json_schema).as("json"))
df = df.select(col("json.*"))

Driver UI for all runs

0 Answers
Related