【问题标题】:Why does Spark Structured Streaming not allow changing the number of input sources?为什么 Spark Structured Streaming 不允许更改输入源的数量?
【发布时间】:2020-09-28 03:16:00
【问题描述】:

我想构建一个从多个 Kafka 主题(数量随时间变化)读取的 Spark 流式传输管道。我打算使用Spark Structured Streaming + Kafka Integration Guide 中列出的两个选项之一来停止流式传输作业,添加/删除新主题,并在需要更新流式传输作业中的主题时再次启动该作业:

# Subscribe to multiple topics
df = spark \
  .readStream \
  .format("kafka") \
  .option("kafka.bootstrap.servers", "host1:port1,host2:port2") \
  .option("subscribe", "topic1,topic2") \
  .load()
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")

# Subscribe to a pattern
df = spark \
  .readStream \
  .format("kafka") \
  .option("kafka.bootstrap.servers", "host1:port1,host2:port2") \
  .option("subscribePattern", "topic.*") \
  .load()
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")

经过进一步调查,我注意到Spark Structured Streaming Programming Guide 中的以下几点,并试图了解为什么“不允许”更改输入源的数量:

输入源的数量或类型(即不同的源)的变化:这是不允许的。

“不允许”的定义(也来自Spark Structured Streaming Programming Guide):

术语不允许意味着您不应该进行指定的更改,因为重新启动的查询很可能会因不可预知的错误而失败。 sdf 表示使用 sparkSession.readStream 生成的流式 DataFrame/Dataset。

我的理解是 Spark Structured Streaming 实现了自己的checkpointing mechanism

如果出现故障或故意关闭,您可以恢复之前查询的进度和状态,并从中断处继续。这是使用检查点和预写日志完成的。您可以使用检查点位置配置查询,并且查询会将所有进度信息(即每个触发器中处理的偏移范围)和正在运行的聚合(例如快速示例中的字数)保存到检查点位置。此检查点位置必须是 HDFS 兼容文件系统中的路径,并且可以在启动查询时在 DataStreamWriter 中设置为选项。

有人可以解释为什么“不允许”更改来源的数量吗?我认为这将是检查点机制的好处之一。

【问题讨论】:

    标签: apache-spark pyspark apache-kafka spark-structured-streaming


    【解决方案1】:

    在现有正在运行的模型流作业中添加新输入源的步骤

    1. 停止模型正在运行的当前正在运行的流。
    2. hdfs dfs -get output/checkpoints/offsets /offsets

    目录中将有 3 个文件(因为最后 3 个偏移量是由 spark 保存的)。下面是单个文件的示例格式

    v1

    { "batchWatermarkMs":0,"batchTimestampMs":1578463128395,"conf":{"spark.sql.streaming.stateStore.providerClass":"org.apache.spark.sql.execution.streaming.state.HDFSBackedStateStoreProvider","spark.sql.streaming.flatMapGroupsWithState.stateFormatVersion":"2","spark.sql.streaming.multipleWatermarkPolicy":"min","spark.sql.streaming.aggregation.stateFormatVersion":"2","spark.sql.shuffle.partitions":"200"}}
    { "logOffset":0}
    { "logOffset":0}
    
    • 每个 {"logOffset":batchId} 代表单个输入源。
    • 要添加新的输入源,请在目录中每个文件的末尾添加“-”。

    更新文件示例 v1

    {"batchWatermarkMs":0,"batchTimestampMs":1578463128395,"conf":{"spark.sql.streaming.stateStore.providerClass":"org.apache.spark.sql.execution.streaming.state.HDFSBackedStateStoreProvider","spark.sql.streaming.flatMapGroupsWithState.stateFormatVersion":"2","spark.sql.streaming.multipleWatermarkPolicy":"min","spark.sql.streaming.aggregation.stateFormatVersion":"2","spark.sql.shuffle.partitions":"200"}}
    {"logOffset":0}
    {"logOffset":0}
    
    • 如果您想添加超过 1 个输入源,则添加“-”等于新输入源的数量。
    • hdfs dfs -put -f /offsets output/checkpoints/offsets

    【讨论】:

      【解决方案2】:

      做你想做的最好的方法是在多线程中运行你的 readStreams。 我正在这样做,同时阅读 40 张桌子。为此,我遵循这篇文章: https://cm.engineering/multiple-spark-streaming-jobs-in-a-single-emr-cluster-ca86c28d1411.

      我将简要介绍一下我在阅读后所做的事情,并将我的代码结构与 main 函数、执行器和一个特征挂载到我的 spark 会话中,该会话将与所有作业共享。

      1.我想阅读的主题的两个列表。

      因此,在 Scala 中,我创建了两个列表。第一个列表是我一直想阅读的主题,第二个列表是动态列表,当我停止工作时,我可以添加一些新主题。

      1. 运行作业的模式匹配。

      我有两个不同的工作,一个运行到我一直运行的表和运行到特定主题的动态作业,换句话说,如果我想添加一个新主题并为他创建一个新工作,我在模式匹配中添加了这个工作。在下面的代码中,我想对 Cars 和 Ship 表运行特定作业,并且我放入特定列表中的所有其他表都将运行相同的复制表作业

        var tables = specifcTables ++ dynamicTables
      
        tables.map(table => {
          table._1 match {
            case "CARS" => new CarsJob
            case "SHIPS" => new ShipsReplicationJob
            case _ => new ReplicationJob
      

      之后,我将此模式匹配传递给 createjobs 函数,该函数将实例化每个作业,并将此函数传递给 startFutureTask 函数,该函数会将这些作业中的每一个放在不同的线程中

      startFutureTasks(createJobs(tables))
      

      我希望我有所帮助。谢谢!

      【讨论】:

        猜你喜欢
        • 2018-02-28
        • 1970-01-01
        • 2019-07-30
        • 1970-01-01
        • 2011-08-15
        • 1970-01-01
        • 2022-01-08
        • 2020-03-19
        • 2020-03-04
        相关资源
        最近更新 更多