如何在Spark Streaming中使用无限的Scala流作为源?

Bas*_*ato 5 scala apache-spark spark-streaming

假设我基本上是想Stream.from(0)作为InputDStream.我该怎么做?我能看到的唯一方法是使用StreamingContext#queueStream,但是我必须从另一个线程或子类中排队元素Queue以创建一个行为类似于无限流的队列,这两者都感觉像是一个黑客.

这样做的正确方法是什么?

Eug*_*nev 2

我认为默认情况下它在 Spark 中不可用,但使用 ReceiverInputDStream 很容易实现它。

import org.apache.spark.storage.StorageLevel
import org.apache.spark.streaming.StreamingContext
import org.apache.spark.streaming.dstream.ReceiverInputDStream
import org.apache.spark.streaming.receiver.Receiver

class InfiniteStreamInputDStream[T](
       @transient ssc_ : StreamingContext,
       stream: Stream[T],
       storageLevel: StorageLevel
      ) extends ReceiverInputDStream[T](ssc_)  {

  override def getReceiver(): Receiver[T] = {
    new InfiniteStreamReceiver(stream, storageLevel)
  }
}

class InfiniteStreamReceiver[T](stream: Stream[T], storageLevel: StorageLevel) extends Receiver[T](storageLevel) {

  // Stateful iterator
  private val streamIterator = stream.iterator

  private class ReadAndStore extends Runnable {
    def run(): Unit = {
      while (streamIterator.hasNext) {
        val next = streamIterator.next()
        store(next)
      }
    }
  }

  override def onStart(): Unit = {
    new Thread(new ReadAndStore).run()    
  }

  override def onStop(): Unit = { }
}
Run Code Online (Sandbox Code Playgroud)