【发布时间】:2015-08-06 20:53:25
【问题描述】:
我在 Spark 驱动程序中有一个模块正在侦听 Kafka 队列,并且根据队列的内容,我需要修改广播变量(或闭包)的内容。在此示例中,这可能是一个字符串。
例如,如果字符串“change”到达队列,我需要更新每个节点中的广播变量。
我希望看到一种干净且高效的模式来执行此操作,或者至少接收关于我可以在哪里找到一些材料以更好地了解如何在 Spark 集群中传播修改的输入。
【问题讨论】:
-
您能否提供一些示例代码来说明您想要做什么?
-
特别是“我在 Spark 驱动程序中有一个模块正在监听 Kafka 队列”令人费解。有关详细信息,请参阅@Bacon 答案的讨论。
标签: scala apache-spark spark-streaming