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.