【问题标题】:Trying to understand how Spark Streaming works?试图了解 Spark Streaming 的工作原理?
【发布时间】: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流中计数,这样走的路对吗?

【问题讨论】:

    标签: python spark-streaming


    【解决方案1】:

    驱动程序是否将定义的转换/计算发送给所有 executors only ONCE AND FOR ALL?

    不,每个任务都被序列化并在每个批次迭代中发送给所有工作人员。想想当你有一个在转换中使用的类的实例时会发生什么,Spark 必须能够将同一个实例及其所有状态发送给每个执行器以进行操作。

    如果是这样,我们就没有机会重新定义/改变 计算,对吧?

    转换定义中的逻辑是不变的,但这并不意味着您不能查询存储影响转换内部数据的信息的第三方。

    例如,假设您有一些外部来源,表明您应该过滤哪些字母。然后,您可以在 DStream 上调用 transform 以从驱动程序获取有关过滤哪个字母的数据。

    或者驱动程序是否在每次执行后重新启动/重新调度所有执行程序 批量数据完成了吗?

    它不会重新启动,它只是在每个批处理间隔开始一个新作业。如果您将StreamingContext 批处理持续时间定义为 60 秒,则每 60 秒一个新作业(微批处理)将开始处理数据。

    根据您的跟进,是的,我就是这样做的。

    【讨论】:

    • 非常感谢您的回复,您对我的第一个问题的回答在我脑海中引发了另一个问题:如果我确实在转换中使用了一个类的实例,那么这个类可以是用户定义的或者来自标准的python/scala/java lib,所以这个特定的实例应该被复制并发送给所有的执行者,对吧?而且要做到这一点,这个实例必须有序列化和反序列化的方法,否则不能发货,也就是不能在transform里面使用,对吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-07-12
    • 2017-09-20
    • 2018-06-23
    • 1970-01-01
    相关资源
    最近更新 更多