【问题标题】:Spark Saving Results to HDFSSpark 将结果保存到 HDFS
【发布时间】:2016-04-27 10:14:56
【问题描述】:

在 HDFS 上使用带有一个主节点和三个工作节点的 spark 保存解析的 XML 文件时,我遇到了奇怪的行为,问题是

当我解析 XMLFile 并尝试保存在 HDFS 中时,文件无法保存所有解析结果。

当我通过指定

使用本地模式执行相同的代码时
sc = SparkContext("local", "parser") 

and the spark-submit will be ./bin/spark-submit xml_parser.py

此运行在 hdfs 上提供了 117mb 的已解析文件,具有完整的记录。

如果在 spark-client 模式下执行代码,那么我做了以下操作,

 sc = SparkContext("spark://master:7077", "parser") 

火花提交是,

./bin/spark-submit --master yarn-client --deploy-mode client --driver-memory 7g --executor-memory 4g  --executor-cores 2  xml_parser.py 1000

给我 19mb 的 hdfs 文件,但记录不完整。

为了在两种情况下保存结果,我都使用 rdd.saveAsTextFile("hdfs://")

我正在使用 spark1.6.1-hadoop2.6 和 Apache hadoop 2.7.2

谁能帮帮我。我不明白为什么会这样。 我有以下 sparkCluster,

1-master 8GbRAM

2-workerNode1 8GbRAM

3-WorkerNode2 8GbRAM

4-workerNode3 8GbRAM

我已经在 Hadoop-2.7.2 上配置了上面的集群,有 1 个主节点和 3 个 DataNode,

如果我 jps On severNode 给了我,

24097 大师

21652 日/秒

23398 名称节点

23799 资源管理器

23630 次要名称节点

所有数据节点上的 JPS,

8006 工人

7819 节点管理器

27164 日元

7678 数据节点

通过检查 HadoopNameNode ui master:9000 给我三个实时数据节点

通过检查 master:7077 上的 SparkMaster Ui 给了我三个 live worker

请看这里,

sc = SpakContext("spark://master:7077", "parser")
--------------------------------------------
 contains the logic of XMLParsing
--------------------------------------------
and I am appending the result in one list like,
cc_list.append([final_cll, Date,Time,int(cont[i]), float(values[i]),0])
Now I am Parallelizing the above cc_list like
 parallel_list = sc.parallelize(cc_list)
 parallel_list.saveAsTextFile("hdfs://master:9000/ some path")
 Now I am Doing some operations here.
 new_list = sc.textFile("hdfs://localhost:9000/some path/part-00000).map(lambda line:line.split(','))

 result = new_list.map(lambda x: (x[0]+',   '+x[3],float(x[4]))).sortByKey('true').coalesce(1)
 result = result.map(lambda x:x[0]+','+str(x[1]))
 result = result.map(lambda x: x.lstrip('[').rstrip(']').replace(' ','')).saveAsTextFile("hdfs://master:9000/some path1))

【问题讨论】:

  • 可以分享代码吗?否则很难理解到底发生了什么......
  • 整个解析逻辑在python中没有火花转换和动作我只使用了我并行化了那个列表。

标签: apache-spark hdfs pyspark


【解决方案1】:

对不起,这里有这样愚蠢的问题。其实我发现了两个问题

1) 在多个工作人员上运行时,

 parallel_list = sc.parallelize(cc_list) 

创建 4-5 个部分文件,parallel_list 保存在 Hdfs 中,part-00000 到 part-00004,在加载 parallel_list 时,您可以在代码中看到上述内容

new_list = sc.textFile(pathto parallel_list/part-00000) ==> so it was taking only the first part.

2) 在 localMode 上运行时,

 parallel_list = sc.parallelize(cc_list) was creating only one part file so i was able to pick whole file at one stroke.

因此,在与工人一起在 spark 上运行时,我想出了两个解决方案

1) 我刚刚在从 parallel_list 创建 new_list 时添加了 part-*

2) 通过使用 spark submit 传递 --conf spark.akka.frameSize=1000 将 spark.akka.frameSize 增加到 10000。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2016-10-17
    • 2015-03-03
    • 2018-06-22
    • 2016-01-29
    • 1970-01-01
    • 2016-10-12
    • 2014-08-21
    • 2019-12-05
    相关资源
    最近更新 更多