【发布时间】: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