【问题标题】:First Pyspark Program第一个 Pyspark 程序
【发布时间】:2020-06-25 09:38:45
【问题描述】:

我在运行我的第一个 Pyspark 程序时遇到了问题。 我在 jypyter 笔记本上运行此代码,我将其配置为使用而不是 shell

import sys
from pyspark import SparkContext

lines = sc.textFile(sys.argv[1])
word_counts = lines.flatMap(lambda line: line.split(' '))\
                   .map(lambda word: (word,1)) \
                   .reduceByKey(lambda count1, count2: count1 + count2) \
                   .collect()

for (word,count) in woord_counts:
    print(word,count)

我收到了这个错误:


Py4JJavaError                             Traceback (most recent call last)
<ipython-input-7-727078dac5d6> in <module>()
      5 sc
      6 lines = sc.textFile(sys.argv[1])
----> 7 word_counts = lines.flatMap(lambda line: line.split(' '))                   .map(lambda word: (word,1))                    .reduceByKey(lambda count1, count2: count1 + count2)                    .collect()
      8 
      9 for (word,count) in word_counts:

/home/mouad/code/spark/python/pyspark/rdd.py in reduceByKey(self, func, numPartitions, partitionFunc)
   1696         [('a', 2), ('b', 1)]
   1697         """
-> 1698         return self.combineByKey(lambda x: x, func, func, numPartitions, partitionFunc)
   1699 
   1700     def reduceByKeyLocally(self, func):

/home/mouad/code/spark/python/pyspark/rdd.py in combineByKey(self, createCombiner, mergeValue, mergeCombiners, numPartitions, partitionFunc)
   1923         """
   1924         if numPartitions is None:
-> 1925             numPartitions = self._defaultReducePartitions()
   1926 
   1927         serializer = self.ctx.serializer

不能粘贴整个错误代码,所以我会在接下来的 cmets 中这样做..

我尝试通过 spark-submit 运行此脚本,但出现以下错误:


hduser_@Master:/home/mouad/code/spark/bin$ ./spark-submit ./wordcount.py ./test.txt
20/03/13 10:26:50 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
Error executing Jupyter command '/home/mouad/code/spark/bin/./wordcount.py': [Errno 2] No such file or directory
log4j:WARN No appenders could be found for logger (org.apache.spark.util.ShutdownHookManager).
log4j:WARN Please initialize the log4j system properly.
log4j:WARN See http://logging.apache.org/log4j/1.2/faq.html#noconfig for more info.

谁能帮我解决这个问题?

【问题讨论】:

  • 错误代码的第二部分` /home/mouad/code/spark/python/pyspark/rdd.py in _defaultReducePartitions(self) 2333 return self.ctx.defaultParallelism 2334 else: -> 2335 return self.getNumPartitions() 2336 2337 def lookup(self, key): /home/mouad/code/spark/python/pyspark/rdd.py in getNumPartitions(self) 2599 2600 def getNumPartitions(self): -> 2601 return self. _prev_jrdd.partitions().size() 2602 2603 @property `
  • 第三部分错误代码/home/mouad/code/spark/python/lib/py4j-0.10.8.1-src.zip/py4j/java_gateway.py in __call__(self, *args) 1284 answer = self.gateway_client.send_command(command) 1285 return_value = get_return_value( -&gt; 1286 answer, self.gateway_client, self.target_id, self.name) 1287 1288 for temp_arg in temp_args:
  • 第四部分:/home/mouad/code/spark/python/pyspark/sql/utils.py in deco(*a, **kw) 96 def deco(*a, **kw): 97 try: ---&gt; 98 return f(*a, **kw) 99 except py4j.protocol.Py4JJavaError as e: 100 converted = convert_exception(e.java_exception)
  • 第五部分:/home/mouad/code/spark/python/lib/py4j-0.10.8.1-src.zip/py4j/protocol.py in get_return_value(answer, gateway_client, target_id, name) 326 raise Py4JJavaError( 327 "调用 {0}{1}{2}.\n" 时发生错误。--> 328 format(target_id, ".", name), value) 329 else: 330 raise Py4JError(
  • 第 6 部分:Py4JJavaError:调用 o35.partitions 时发生错误。 :org.apache.hadoop.mapred.InvalidInputException:输入路径不存在:文件:/home/mouad/code/spark/-f at org.apache.hadoop.mapred.FileInputFormat.singleThreadedListStatus(FileInputFormat.java:297) at org.apache.hadoop.mapred.FileInputFormat.listStatus(FileInputFormat.java:239) at org.apache.hadoop.mapred.FileInputFormat.getSplits(FileInputFormat.java:325) at org.apache.spark.rdd.HadoopRDD.getPartitions( HadoopRDD.scala:205) 在 org.apache.spark.rdd.RDD.$anonfun$partitions$2(RDD.scala:27​​6)

标签: python apache-spark pyspark log4j


【解决方案1】:

如果您已将 Jupyter Notebook 配置为运行 Pyspark 应用程序,则无需在代码中使用 sys.argv[1] 来传递文件的位置,而是应该传递文件的完整路径。

如果你的spark master和slave节点运行正常,可以尝试执行以下代码。

from operator import add
from pyspark.sql import SQLContext, SparkSession

spark = (
    SparkSession.builder.appName("WordCountApp")
    .master("local[*]")
    .getOrCreate()
)
sqlcontext = SQLContext(spark)

lines  = spark.read.text("wordcount.txt").rdd.map(lambda r: r[0])
counts = lines.flatMap(lambda x: x.split(' ')) \
              .map(lambda x: (x, 1)) \
              .reduceByKey(add)

output = counts.collect()
for (word, count) in output:
    print("%s: %i" % (word, count))

【讨论】:

    【解决方案2】:

    我注意到您使用 ./test.txt 将文件路径作为相对路径,请使用 file:///&lt;full-path-of-file&gt; 表示法提供完整文件路径,您也在集群上运行它,请确保您尝试使用的文件集群中所有节点的同一位置都可以读取。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-12-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多