【问题标题】:What is LocalTableScan in Spark Structured Streaming for?Spark 结构化流中的 LocalTableScan 有什么用?
【发布时间】:2019-02-13 21:58:14
【问题描述】:

有谁知道 Spark Structured Streaming 中 LocalTableScan 对应的是什么?

我试图了解在本地 [*] 模式下运行的 Spark 结构流应用程序中观察到的一种奇怪行为。

我的机器上有 8 个内核。虽然我的大多数批次都有 8 个分区,但每隔一段时间我就会得到 16 个或 32 个或 56 个分区/任务等等。我注意到它始终是 8 的倍数。我在打开阶段选项卡时注意到,当它发生时,是因为有多个 LocalTableScan。

也就是说,如果我有 2 个 LocalTableScan,那么小批量作业将有 16 个任务/分区等等。

为了提供一些上下文,因为我怀疑它可能来自它,我正在使用 MemoryStream。

val rows = MemoryStream[Map[String,String]]
val df = rows.toDF()
val rdf = df.mapPartitions{ it => {.....}}(RowEncoder.apply(StructType(List(StructField("blob", StringType, false)))))

我有一个未来会像这样喂我的记忆流:

Future {
    blocking {
      for (i <- 1 to 100000) {
        rows.addData(maps)
        Thread.sleep(3000)
      }
    }
  }

然后是我的查询:

rdf.writeStream.
    trigger(Trigger.ProcessingTime("1 seconds"))
    .format("console").outputMode("append")
    .queryName("SourceConvertor1").start().awaitTermination()

请问有什么建议吗?提示?

【问题讨论】:

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


    【解决方案1】:

    它表示在驱动程序的内存中。如您的代码所示。

    【讨论】:

    • 嗨@thebluephantom 非常感谢您花时间提供帮助。同时,请您进一步解释您的意思:“它在驱动程序的内存中表示”?您的意思是说,spark 正在扫描位于驱动程序上的内存表中的数据吗?知道为什么,对于一个小批量,我得到 2 次扫描吗? 2 取?换句话说,一个小批量似乎对应于从 MemoryStream 获取/扫描的 2 个批次。这对你有意义吗?
    • 我回答了问题的标题。您需要提出一个新问题。
    • 从我收集的信息来看,在我看来,作业计划可以说是 32 个任务/分区,但由于可用的核心数量,当时有效地执行了 8 个任务。但是,我看不出与内存表中的扫描/获取次数有什么关系?当它每秒要做一个小批量作业时,为什么还要进行多次扫描?我只是不明白。我的意思是链接很奇怪,因为它可以很好地进行 2 次扫描,但比将其作为一批提供给作业,因此有很大的分区。
    • 在 hols 上,以后看起来会很有趣。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-05-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-21
    • 1970-01-01
    相关资源
    最近更新 更多