Apache Spark - Exception handling inside foreachRDD

Viewed 96

guys. I´m runnig a Spark Streaming job that reads from Kafka, do some things and put the results in another Kafka topic (Spark 2.1.1)

I notice when one task throws an exception, the execution of foreachRDD finalizes "normally", committing the offsets that I handle manually, what according to my understanding, it was not expected.

The code is like

kafkaStream.foreachRDD { rdd =>
        if (!rdd.isEmpty()) {
          val offsets = offsetStore.getOffsetsFromRDD(rdd)        
          doTheThing(rdd)          
          offsetStore.persistOffsets(offsets)          
        }
 }

doTheThing() executes a foreachPartition and one of this tasks fail and throws an exception, like this:

def doTheThing(rdd: RDD) = {
    rdd.foreachPartition {l => 
        try {
            doAnotherThing()
        } catch {
             case e: NeedToAbortSomeException => throw e
        }
    }
}

What is wrong? According Spark documentation (http://spark.apache.org/docs/latest/streaming-kafka-0-10-integration.html), is ok to commit a transaction inside foreachRDD block.

0 Answers
Related