Spark streaming: HBase connection is closed when use hbaseMapPartitions

Viewed 235

in my Spark streaming application i use a HBaseContext to put some values into HBase, one put operation for each processed message.

If i use hbaseForeachPartitions, everything is ok.

 dStream
  .hbaseForeachPartition(
    hbaseContext,
    (iterator, connection) => {
      val table = connection.getTable("namespace:table")
      // putHBase is external function in the same Scala object
      val results = iterator.flatMap(packet => putHBaseAndOther(packet))
      table.close()
      results
    }
 )

Instead with hbaseMapPartitions the connection to HBase is closed.

 dStream
  .hbaseMapPartition(
    hbaseContext,
    (iterator, connection) => {
      val table = connection.getTable("namespace:table")
      // putHBase is external function in the same Scala object
      val results = iterator.flatMap(packet => putHBaseAndOther(packet))
      table.close()
      results
    }
 )

Someone can explain me why?

0 Answers
Related