spark-dev mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From badgerpants <>
Subject practical usage of the new "exactly-once" supporting DirectKafkaInputDStream
Date Thu, 30 Apr 2015 17:10:24 GMT
We're a group of experienced backend developers who are fairly new to Spark
Streaming (and Scala) and very interested in using the new (in 1.3)
DirectKafkaInputDStream impl as part of the metrics reporting service we're

Our flow involves reading in metric events, lightly modifying some of the
data values, and then creating aggregates via reduceByKey. We're following
the approach in Cody Koeninger's blog on exactly-once streaming
( in
which the Kakfa OffsetRanges are grabbed from the RDD and persisted to a
tracking table within the same db transaction as the data within said

Within a short time frame the offsets in the table fall out of synch with
the offsets. It appears that the writeOffsets method (see code below)
occasionally doesn't get called which also indicates that some blocks of
data aren't being processed either; the aggregate operation makes this
difficult to eyeball from the data that's written to the db.

Note that we do understand that the reduce operation alters that
size/boundaries of the partitions we end up processing. Indeed, without the
reduceByKey operation our code seems to work perfectly. But without the
reduceByKey operation the db has to perform *a lot* more updates. It's
certainly a significant restriction to place on what is such a promising
approach. I'm hoping there simply something we're missing.

Any workarounds or thoughts are welcome. Here's the code we've got:

def run(stream: DStream[Input], conf: Config, args: List[String]): Unit = {
    val sumFunc: (BigDecimal, BigDecimal) => BigDecimal = (_ + _)

    val transformStream = stream.transform { rdd =>
      val offsets = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
      printOffsets(offsets) // just prints out the offsets for reference
      rdd.mapPartitionsWithIndex { case (i, iter) =>
        iter.flatMap { case (name, msg) => extractMetrics(msg) }
          .map { case (k,v) => ( ( keyWithFlooredTimestamp(k), offsets(i) ),
v ) }
    }.reduceByKey(sumFunc, 1)

    transformStream.foreachRDD { rdd =>
      rdd.foreachPartition { partition =>
        val conn = DriverManager.getConnection(dbUrl, dbUser, dbPass)
        val db = DB(conn)

        db.autoCommit { implicit session =>
          var currentOffset: OffsetRange = null
          partition.foreach { case (key, value) =>
            currentOffset = key._2
            writeMetrics(key._1, value, table)
          writeOffset(currentOffset) // updates the offset positions


View this message in context:
Sent from the Apache Spark Developers List mailing list archive at

To unsubscribe, e-mail:
For additional commands, e-mail:

View raw message