【发布时间】:2018-12-07 11:54:20
【问题描述】:
我有一个从 Kafka 读取数据的 spark 流应用程序(使用 spark 1.6.1)。我正在使用 hdfs 上的 spark 检查点目录从故障中恢复。
代码在每批开始时使用带有映射的广播变量,如下所示
public static void execute(JavaPairDStream<String, MyEvent> events) {
final Broadcast<Map<String, String>> appConfig = MyStreamConsumer.ApplicationConfig.getInstance(new JavaSparkContext(events.context().sparkContext()));
当我在运行时提交作业时,一切正常,来自 kafka 的所有事件都得到正确处理。从故障中恢复时会出现问题(通过重新启动运行 spark 的机器进行测试) - spark 流应用程序实际上正确启动,并且在 UI 中一切看起来都很好,作业正在运行,但是一旦通过以下异常发送数据调用此方法时遇到(&作业崩溃):
appConfig.value() (the broadcast variable from the start!)
由于 Spark 作业错误而失败
Caused by: java.lang.ClassCastException: org.apache.spark.util.SerializableConfiguration cannot be cast to java.util.Map
如果我在 spark UI 中终止驱动程序并从命令行重新提交作业,一切都会再次正常运行。 但它对我们产品的要求是它可以从故障中自动恢复,甚至只是重新启动任何集群节点,所以我必须修复上述问题。 这个问题肯定与使用 Broadcast 变量和在重启后从 spark checkpoint 目录加载状态有关
另外请注意,我确实正确地创建了广播实例(懒惰/单例):
public static Broadcast<Map<String, String>> getInstance(JavaSparkContext sparkContext) {
if (instance == null) {
我确实意识到这个问题似乎与: Is it possible to recover an broadcast value from Spark-streaming checkpoint
但我无法按照他们的说明解决问题
【问题讨论】:
标签: apache-spark