【问题标题】:pyspark - kafka stream - out of memorypyspark - kafka 流 - 内存不足
【发布时间】:2019-05-13 12:37:10
【问题描述】:

我正在尝试使用此代码测试带有代理版本 0.10 的 kafka 流。这只是打印主题内容的简单代码。还没有什么大不了的!但是,由于某种原因,内存不够用(VM 中有 10GB 的 RAM)!代码:

# coding: utf-8

"""
kafka-test-003.py: test with broker 0.10(new Spark Stream API)

How to run this script?

spark-submit --jars jars/spark-sql-kafka-0-10_2.11-2.3.0.jar,jars/kafka-clients-0.11.0.0.jar kafka-test-003.py



"""


import pyspark 
from pyspark import SparkContext
from pyspark.sql.session import SparkSession,Row
from pyspark.sql.types import *
from pyspark.sql.functions import *
from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils


# starting spark session
spark = SparkSession.builder.appName("Kakfa-test").getOrCreate()
spark.sparkContext.setLogLevel('WARN')

# getting streaming context
sc = spark.sparkContext
ssc = StreamingContext(sc, 2) # batching duration: each 2 seconds

broker = "kafka.some.address:9092"
topic = "my.topic"

### Streaming

df = spark \
  .readStream \
  .format("kafka") \
  .option("kafka.bootstrap.servers", broker) \
  .option("startingOffsets", "earliest") \
  .option("subscribe", topic) \
  .load() \
  .select(col('key').cast("string"),col('value').cast("string"))

query = df \
  .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \
  .writeStream \
  .outputMode("append") \
  .format("console") \
  .start()

### End Streaming

query.awaitTermination()

运行火花提交:

spark-submit --master local[*] --driver-memory 5G --executor-memory 5G --jars jars/kafka-clients-0.11.0.0.jar,jars/spark-sql-kafka-0-10_2.11-2.3.0.jar kafka-test-003.py

不幸的是,结果是:

java.lang.OutOfMemoryError: Java 堆空间

我假设 Kafka 应该每次带来一小部分数据来避免这个问题,对吧?那么,我做错了什么?

【问题讨论】:

    标签: pyspark apache-kafka out-of-memory


    【解决方案1】:

    spark 内存管理是一个复杂的过程。最佳解决方案不仅取决于您的数据和操作类型以及系统行为 您可以重试以下 spark 命令吗:

    spark-submit --master local[*] --driver-memory 4G --executor-memory 2G --executor-cores 5 --num-executors 8 --jars jars/kafka-clients-0.11.0.0.jar,jars/spark-sql-kafka-0-10_2.11-2.3.0.jar kafka-test-003.py

    您可以通过调整性能来按照以下链接调整上述内存参数吗? Using spark-submit, what is the behavior of the --total-executor-cores option?

    【讨论】:

      猜你喜欢
      • 2015-12-02
      • 2018-07-11
      • 2016-10-10
      • 1970-01-01
      • 1970-01-01
      • 2019-03-11
      • 2019-07-31
      • 1970-01-01
      • 2020-05-17
      相关资源
      最近更新 更多