【发布时间】:2019-10-09 18:31:31
【问题描述】:
Spark 结构化流执行器因 OutOfMemoryError 而失败
使用 VirtualVM 检查堆分配表明 JMX Mbean Server 内存使用量随时间线性增长。
经过进一步调查,似乎 JMX Mbean 充满了数千个 KafkaMbean 对象实例,其中消费者指标 (\d+) 达到数千个(等于在执行程序上创建的任务数)。
在执行器上运行带有 DEBUG 日志的 Kafka 消费者表明,执行器添加了数千个指标传感器,并且通常根本不删除它们或只删除一些
我正在运行 HDP Spark 2.3.0.2.6.5.0-292 和 HDP Kafka 1.0.0.2.6.5.0-292。
这是我初始化结构化流的方法:
sparkSession
.readStream
.format("kafka")
.options(Map("kafka.bootstrap.servers" -> KAFKA_BROKERS,
"subscribePattern" -> INPUT_TOPIC,
"startingOffsets" -> "earliest",
"failOnDataLoss" -> "false"))
.mapPartitions(processData)
.writeStream
.format("kafka")
.options(Map("kafka.bootstrap.servers" -> KAFKA_BROKERS,
"checkpointLocation" -> CHECKPOINT_LOCATION))
.queryName("Process Data")
.outputMode("update")
.trigger(Trigger.ProcessingTime(1000))
.load()
.start()
.awaitTermination()
我期待 Spark/Kafka 在任务完成时正确清理 MBean,但事实并非如此。
【问题讨论】:
标签: apache-spark apache-kafka spark-structured-streaming