如何在单个Spark应用中创建多个StreamingContext且无需启用多上下文允许
解决Spark-Shell中创建StreamingContext的多上下文冲突问题
这个问题我在使用spark-shell做Streaming测试时也碰到过,核心原因很简单:spark-shell启动时已经自动帮你初始化了一个SparkContext实例(就是默认的sc变量),当你自己通过SparkConf新建StreamingContext时,它会默认创建一个全新的SparkContext,这就触发了Spark不允许多上下文共存的限制。
不用开启spark.driver.allowMultipleContexts这个不推荐的配置,有两种更规范的解决方式:
方式一:复用Spark-Shell已有的SparkContext
直接用shell预创建好的sc来初始化StreamingContext,完全避免新上下文的创建:
import org.apache.spark.streaming._ import org.apache.spark.streaming.StreamingContext._ // 复用spark-shell自带的sc实例 val ssc = new StreamingContext(sc, Seconds(30))
这样既符合Spark的设计规范,也不会出现任何警告或错误。
方式二:启动Spark-Shell时指定自定义配置
如果你的Streaming任务需要特定的Spark配置参数,可以在启动spark-shell时通过命令行参数传递,这样预创建的sc就会带上这些配置,之后再用它创建StreamingContext即可:
spark-shell --master local[2] --appName NetworkWordCount --conf spark.some.custom.param=value
进入shell后,直接用方式一的代码创建ssc就行,所有配置都会继承自预初始化的sc。
补充说明
在生产环境的独立Spark应用中,你可以完全控制SparkContext的生命周期,不会出现这种问题;但spark-shell作为交互式环境,为了简化操作提前创建了SparkContext,这才导致了冲突。所以在shell里做Streaming测试时,一定要优先复用已有的sc实例。
内容的提问来源于stack exchange,提问作者shanlodh
相关产品推荐
相关产品推荐

