【发布时间】:2017-06-12 06:19:02
【问题描述】:
我的要求是处理股票市场的每小时数据。 即,每个流式传输间隔从源获取一次数据并通过 DStream 进行处理。
我已经实现了一个自定义接收器,通过实现 onStart() 和 onStop() 方法及其工作来废弃/监控网站。
遇到的挑战:
- 接收线程连续获取数据,即每个间隔多次获取数据。
- 无法协调接收器和 DStream 执行时间间隔。
我尝试过的选项:
- 接收线程休眠几秒钟(等于流传输间隔)。 在这种情况下,数据不是处理时的最新数据。
class CustomReceiver(interval: Int)
extends Receiver[String](StorageLevel.MEMORY_AND_DISK_2) {
def onStart() {
new Thread("Website Scrapper") {
override def run() { receive() }
}.start()
}
def onStop() {
}
/** Create a socket connection and receive data until receiver is stopped */
private def receive() {
println("Entering receive:" + new Date());
try {
while (!isStopped) {
val scriptsLTP = StockMarket.getLiveStockData()
for ((script, ltp) <- scriptsLTP) {
store(script + "," + ltp)
}
println("sent data")
System.out.println("going to sleep:" + new Date());
Thread.sleep(3600 * 1000);
System.out.println("awaken from sleep:" + new Date());
}
println("Stopped receiving")
restart("Trying to connect again")
} catch {
case t: Throwable =>
restart("Error receiving data", t)
}
println("Exiting receive:" + new Date());
}
}
如何使 Spark Streaming 接收器与 DStream 处理同步?
【问题讨论】:
-
是否可以在流式传输间隔开始时获取数据?
标签: apache-spark spark-streaming