【问题标题】:Spark Structured Streaming throwing Java OOM immediatelySpark Structured Streaming 立即抛出 Java OOM
【发布时间】:2018-05-07 22:11:28
【问题描述】:

我正在尝试构建一个简单的管道,使用 Kafka 作为 Spark 结构化流 API 的流源,执行分组聚合并将结果保存到 HDFS。

但是,一旦我提交作业,即使流数据量非常少,我也会收到 Java 堆空间错误。

下面是pyspark中的代码:

allEvents =spark \
    .readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe","MyNewTopic") \
    .option("group.id","aggStream") \
    .option("startingOffsets", "earliest") \
    .load() \
    .select(col("value").cast("string"))

aaIDF = allEvents.filter(col("value").contains("myNewAPI")).select(from_json(col("value"),aaISchema) \
 .alias("colName")).select(col("colName.eventTime"), col("colName.appId"),col("colName.articleId"),col("colName.locale"),col("colName.impression"))

windowedCountsDF = aaIDF.withWatermark("eventTime","10 minutes") \
    .groupBy("appId","articleId","locale",window("eventTime", "2 minutes")).sum("impression").withColumnRenamed("sum(impression)", "views")


query = windowedCountsDF \
    .writeStream \
    .outputMode("append") \
    .format("parquet") \
    .option("path", "/CDS/events/JS/agg/" + strftime("%Y/%m/%d/%H/%M", gmtime()) + "/") \
    .option("checkpointLocation", "/CDS/checkpoint/").start()

以下是例外:

17/11/23 14:24:45 ERROR Utils: Aborting task
java.lang.OutOfMemoryError: Java heap space
    at org.apache.spark.sql.catalyst.expressions.codegen.BufferHolder.grow(BufferHolder.java:73)
    at org.apache.spark.sql.catalyst.expressions.codegen.UnsafeRowWriter.write(UnsafeRowWriter.java:214)
    at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIterator.agg_doAggregateWithKeys$(Unknown Source)
    at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIterator.processNext(Unknown Source)
    at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
    at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$8$$anon$1.hasNext(WholeStageCodegenExec.scala:395)
    at org.apache.spark.sql.execution.datasources.FileFormatWriter$SingleDirectoryWriteTask.execute(FileFormatWriter.scala:315)
    at org.apache.spark.sql.execution.datasources.FileFormatWriter$$anonfun$org$apache$spark$sql$execution$datasources$FileFormatWriter$$executeTask$3.apply(FileFormatWriter.scala:258)
    at org.apache.spark.sql.execution.datasources.FileFormatWriter$$anonfun$org$apache$spark$sql$execution$datasources$FileFormatWriter$$executeTask$3.apply(FileFormatWriter.scala:256)
    at org.apache.spark.util.Utils$.tryWithSafeFinallyAndFailureCallbacks(Utils.scala:1375)
    at org.apache.spark.sql.execution.datasources.FileFormatWriter$.org$apache$spark$sql$execution$datasources$FileFormatWriter$$executeTask(FileFormatWriter.scala:261)
    at org.apache.spark.sql.execution.datasources.FileFormatWriter$$anonfun$write$1$$anonfun$apply$mcV$sp$1.apply(FileFormatWriter.scala:191)
    at org.apache.spark.sql.execution.datasources.FileFormatWriter$$anonfun$write$1$$anonfun$apply$mcV$sp$1.apply(FileFormatWriter.scala:190)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
    at org.apache.spark.scheduler.Task.run(Task.scala:108)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:335)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
    at java.lang.Thread.run(Thread.java:748)

【问题讨论】:

  • 如何提交作业?什么是spark-submit 和选项?
  • @JacekLaskowski spark-submit --packages 'org.mongodb.spark:mongo-spark-connector_2.11:2.2.0,org.apache.spark:spark-sql-kafka-0-10_2 .11:2.1.0,org.apache.spark:spark-streaming-kafka-0-8_2.11:2.1.1' ./AggEventStructuredStreamListener.py
  • 您能否查看jconsole 并查看应用程序的内存要求(并相应地进行调整)?我看不出它可能失败的任何明显原因。请注意,--packages 选项包括用于不同 Spark 版本的库 - 2.2.0 和 2.1.0‌。另外,您正在使用spark-sql-k‌​afkasp‌​ark-streaming-kafka,我怀疑您是否真的需要。摆脱org.apache.spark:sp‌​ark-streaming-kafka-‌​0-8_2.11:2.1.1

标签: apache-spark spark-streaming databricks


【解决方案1】:

两个可能的原因:

  1. 您的水印设置没有生效。您应该使用colName.eventTime 引用该列。

    由于未定义水印(仅在其他类别中定义),因此不会丢弃旧的聚合状态。

  2. 对于 Spark,您应该为 --driver-memory--executor-memory 设置更大的值。

【讨论】:

    【解决方案2】:

    在提交作业时,您需要有适当的驱动程序并执行内存集。这个post 让您简要了解如何设置这些配置。

    【讨论】:

      猜你喜欢
      • 2020-09-08
      • 2020-03-19
      • 2020-09-12
      • 1970-01-01
      • 2022-01-05
      • 2018-06-04
      • 2022-10-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多