【问题标题】:Spark Structured Streaming custom StateStoreProvideSpark 结构化流自定义 StateStoreProvide
【发布时间】:2018-12-08 15:05:40
【问题描述】:

默认情况下,结构化流作业使用HDFSStateStoreProvide。使用 HDFS 存储的问题是它不可扩展。 当作业在高流量时段从 kafka 获取更多数据时,由于以下错误而失败:

18/12/06 15:54:35 ERROR scheduler.TaskSetManager: Task 191 in stage 231.0 failed 4 times; aborting job
18/12/06 15:54:35 ERROR streaming.StreamExecution: Query eventQuery [id = 42051afe-b1bc-438d-8143-2d7e5def717c, runId = 6201c769-b115-4b92-bad5-450b8803b88b] terminated with error
org.apache.spark.SparkException: Job aborted due to stage failure: Task 191 in stage 231.0 failed 4 times, most recent failure: Lost task 191.3 in stage 231.0 (TID 24016, sparkstreamingc1n5.host.bo1.csnzoo.com, executor 659): java.io.EOFException
    at java.io.DataInputStream.readInt(DataInputStream.java:392)
    at org.apache.spark.sql.execution.streaming.state.HDFSBackedStateStoreProvider.org$apache$spark$sql$execution$streaming$state$HDFSBackedStateStoreProvider$$readSnapshotFile(HDFSBackedStateStoreProvider.scala:481)
    at org.apache.spark.sql.execution.streaming.state.HDFSBackedStateStoreProvider$$anonfun$org$apache$spark$sql$execution$streaming$state$HDFSBackedStateStoreProvider$$loadMap$1.apply(HDFSBackedStateStoreProvider.scala:359)
    at org.apache.spark.sql.execution.streaming.state.HDFSBackedStateStoreProvider$$anonfun$org$apache$spark$sql$execution$streaming$state$HDFSBackedStateStoreProvider$$loadMap$1.apply(HDFSBackedStateStoreProvider.scala:358)
    at scala.Option.getOrElse(Option.scala:121)

如何配置自定义状态存储提供?

出于测试目的,我尝试添加一个假类

--conf spark.sql.streaming.stateStore.providerClass=com.streaming.state.RocksDBStateStoreProvider

但是即使这个类不存在,工作仍然是选择 HDFSStateStoreProvider。这是预期的行为吗?

我可以使用任何键值数据库来编写自定义状态提供程序吗?

或者只限于RocksDBCassandra

【问题讨论】:

  • 有 Cassandra StateStore 实现吗?我找不到它。

标签: java apache-spark spark-structured-streaming


【解决方案1】:

如何配置自定义状态存储提供?

您配置自定义状态存储提供程序的方法看起来是正确的,但是一旦您之前运行查询,您就无法更改状态存储提供程序。 (Spark 会从检查点的元数据中读取配置。)这个限制是有意义的,因为在更改状态存储提供程序时,不能保证状态会被恢复。

我可以使用任何键值数据库来编写自定义状态提供程序吗?

一旦您的自定义状态提供程序在状态存储提供程序上实施规范,就没有具体限制。需要考虑的两个主要事项是 1. Spark 将检查每个批次的更改 2. Spark 要求状态存储提供程序恢复特定版本的状态。您的自定义状态提供程序应该是高性能的 - 因为它会增加每批的延迟。

由于自定义状态提供程序,您可能还需要考虑将(传递)依赖项添加到 Spark 应用程序。

【讨论】:

    猜你喜欢
    • 2018-12-20
    • 1970-01-01
    • 2017-05-04
    • 1970-01-01
    • 1970-01-01
    • 2019-04-30
    • 2017-03-06
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多