【发布时间】:2017-03-31 19:48:02
【问题描述】:
这可能是一个愚蠢的问题,但我似乎找不到任何文档用纯英语(好吧,夸张)来澄清这一点,并且在阅读了官方文档和一些博客之后,我仍然对如何驱动程序和执行者工作。
这是我目前的理解:
1) 驱动程序定义转换/计算。
2) 一旦我们调用SparkContext.start(),驱动程序会将定义的转换/计算发送给所有的执行器,这样每个执行器都知道如何处理传入的RDD流数据。
好的,我有一些令人困惑的问题:
1) 驱动程序是否将定义的转换/计算仅发送给所有执行程序ONCE AND FOR ALL?
如果是这样,我们就没有机会重新定义/更改计算,对吧?
例如,我做了一个类似于one 的字数统计工作,但是我的工作有点复杂,我只想计算前 60 秒以字母 J 开头的单词,然后只计算以字母开头的单词在接下来的 60 年代用字母 K,然后只有以 ...... 开头的单词,这样继续下去。
那么我应该如何在驱动程序中实现这个流式传输作业?
2) 还是说驱动在每批数据做完后重启/重新调度所有的executor?
跟进
为了解决1)的问题,我想我可以利用一些外部存储介质,比如redis,我的意思是我可以在驱动程序中实现一个处理函数count_fn,并且每次当这个count_fn被调用,它会从redis中读取来获取起始字母,然后在RDD流中计数,这样走的路对吗?
【问题讨论】: