【问题标题】:Pyspark error for java heap space errorjava堆空间错误的Pyspark错误
【发布时间】:2016-08-24 12:10:32
【问题描述】:

我是使用 Spark 1.6.1两个工人 的新手,每个工人都有 1GB 内存5 个核心 分配,在 33MB 文件上运行此代码。

此代码用于在 spark 中索引单词。

from textblob import TextBlob as tb
from textblob_aptagger import PerceptronTagger
import numpy as np
import nltk.data
import Constants
from pyspark import SparkContext,SparkConf
import nltk

TOKENIZER = nltk.data.load('tokenizers/punkt/english.pickle')

def word_tokenize(x):
   return nltk.word_tokenize(x)

def pos_tag (s):
  global TAGGER
  return TAGGER.tag(s)

def wrap_words (pair):
  ''' associable each word with index '''
  index = pair[0]
  result = []
  for word, tag in pair[1]:
    word = word.lower()
    result.append({ "index": index, "word": word, "tag": tag})
    index += 1
  return result

if __name__ == '__main__':

  conf = SparkConf().setMaster(Constants.MASTER_URL).setAppName(Constants.APP_NAME)
  sc = SparkContext(conf=conf)
  data = sc.textFile(Constants.FILE_PATH)

  sent = data.flatMap(word_tokenize).map(pos_tag).map(lambda x: x[0]).glom()
  num_partition = sent.getNumPartitions()
  base = list(np.cumsum(np.array(sent.map(len).collect())))
  base.insert(0, 0)
  base.pop()
  RDD = sc.parallelize(base,num_partition)
  tagged_doc = RDD.zip(sent).map(wrap_words).cache()

对于小于 25MB 的较小文件,代码可以正常工作,但对于大于 25MB 的文件会出错。
帮我解决此问题或提供此问题的替代方法?

【问题讨论】:

    标签: python numpy optimization pyspark


    【解决方案1】:

    这是因为 .collect()。当您将 rdd 转换为经典的 Python 变量(或 np.array)时,您会失去一切,所有数据都收集在同一个地方。

    【讨论】:

    • 你能建议一个替代解决方案吗?
    • 你会解释你在做什么,这是一个不同的问题。无论如何,如果您使用 spark,collect 仅用于调试,请坚持 rdd 操作。有很多方法可以在 pyspark 中编写累积和,如果您在某个地方遇到困难,请在进行一些研究后打开另一个讨论。
    猜你喜欢
    • 2023-04-07
    • 1970-01-01
    • 2011-10-10
    • 1970-01-01
    • 1970-01-01
    • 2016-01-21
    • 2012-06-05
    • 2018-06-16
    相关资源
    最近更新 更多