【发布时间】:2017-04-08 12:52:49
【问题描述】:
我是 Spark 的新手,刚刚在我的集群上运行它(运行社区版 MapR 的 9 节点集群上的 Spark 2.0.1)。我通过
提交字数示例./bin/spark-submit --master yarn --jars ~/hadoopPERMA/jars/hadoop-lzo-0.4.21-SNAPSHOT.jar examples/src/main/python/wordcount.py ./README.md
并得到以下输出
17/04/07 13:21:34 WARN Client: Neither spark.yarn.jars nor spark.yarn.archive is set, falling back to uploading libraries under SPARK_HOME.
: 68
help: 1
when: 1
Hadoop: 3
...
看起来一切正常。当我添加 --deploy-mode cluster 时,我得到以下输出:
17/04/07 13:23:52 WARN Client: Neither spark.yarn.jars nor spark.yarn.archive is set, falling back to uploading libraries under SPARK_HOME.
所以没有错误但我没有看到字数统计结果。我错过了什么?我在我的历史服务器中看到了该作业,它说它已成功完成。我还检查了 DFS 中的用户目录,但除了这个空目录之外没有写入新文件:/user/myuser/.sparkStaging
代码(Spark 附带的 wordcount.py 示例):
from __future__ import print_function
import sys
from operator import add
from pyspark.sql import SparkSession
if __name__ == "__main__":
if len(sys.argv) != 2:
print("Usage: wordcount <file>", file=sys.stderr)
exit(-1)
spark = SparkSession\
.builder\
.appName("PythonWordCount")\
.getOrCreate()
lines = spark.read.text(sys.argv[1]).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))
spark.stop()
【问题讨论】:
标签: apache-spark pyspark hadoop-yarn