【发布时间】:2020-07-22 22:50:12
【问题描述】:
我对 spark 结构流有一个非常奇怪的问题。 Spark 结构流式处理为每个微批次创建两个 Spark 作业。 结果,从 Kafka 读取数据两次。 这是一个简单的代码sn-p。
import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.streaming.Trigger
object CheckHowSparkReadFromKafka {
def main(args: Array[String]): Unit = {
val session = SparkSession.builder()
.config(new SparkConf()
.setAppName(s"simple read from kafka with repartition")
.setMaster("local[*]")
.set("spark.driver.host", "localhost"))
.getOrCreate()
val testPath = "/tmp/spark-test"
FileSystem.get(session.sparkContext.hadoopConfiguration).delete(new Path(testPath), true)
import session.implicits._
val stream = session
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka-20002-prod:9092")
.option("subscribe", "topic")
.option("maxOffsetsPerTrigger", 1000)
.option("failOnDataLoss", false)
.option("startingOffsets", "latest")
.load()
.repartitionByRange( $"offset")
.writeStream
.option("path", testPath + "/data")
.option("checkpointLocation", testPath + "/checkpoint")
.format("parquet")
.trigger(Trigger.ProcessingTime(10.seconds))
.start()
stream.processAllAvailable()
发生这种情况是因为如果.repartitionByRange( $"offset"),如果我删除此行,一切都很好。
但是使用 spark 创建两个作业,一个具有 1 个阶段刚刚从 Kafka 读取,第二个具有 3 个阶段读取 -> 洗牌 -> 写入。
所以第一份工作的结果从未使用过。
这会对性能产生重大影响。 我的一些 Kafka 主题有 1550 个分区,因此阅读它们两次很重要。 如果我添加缓存,事情会变得更好,但这对我来说不是一种方式。 在本地模式下,批处理中的第一个作业用时不到 0.1 毫秒,但索引为 0 的批处理除外。但在 YARN 集群和 Messos 中,这两个作业完全符合预期,而我的主题则需要将近 1.2 分钟。
为什么会这样?我怎样才能避免这种情况?看起来像虫子?
附:我使用火花 2.4.3。
【问题讨论】:
-
你能看到write阶段正在创建的任务号和shuffle write/read吗?
-
@EmiCareOfCell44 是的,我可以看到所有任务编号和阶段。第一个 Job 有 1 个阶段,其中任务数为 240(与 Kafka 中的分区数相同)。第二个 Job 有 2 个 stage,其中一个 stage 类似于第一个 Job,第二个 stage 是 shuffle + write。
-
我想对于范围分区器,Spark 必须首先读取所有消息才能创建 RangePartiitoner。因为偏移量是每个分区中的一个未知数,它读取并重新分配一个分区中的所有消息以创建de索引,然后对每个执行器要处理的数据进行混洗。您的案例必须使用范围分区器吗?
-
1.是的,RangePartitioner 是强制性的。 2. Spark 不能在另一个作业中重用作业的结果,除非将其持久化到某个接收器并再次读取。并且大多数和第二个工作1阶段的事情和第一个阶段的唯一阶段完全一样。更重要的是,如果我将使用直接流式传输并从每个微批次红色创建一个数据帧,然后使用相同的逻辑 - 我将有 1 个作业。
-
当我使用简单排序时也会发生同样的情况。
标签: scala apache-spark spark-structured-streaming spark-streaming-kafka