如何将生成的RDD写入Spark python中的csv文件

Jas*_*ald 21 python csv file-writing apache-spark pyspark

我有一个结果RDD labelsAndPredictions = testData.map(lambda lp: lp.label).zip(predictions).这有以这种格式输出:

[(0.0, 0.08482142857142858), (0.0, 0.11442786069651742),.....]
Run Code Online (Sandbox Code Playgroud)

我想要的是创建一个CSV文件,其中一列labels(上面输出中的元组的第一部分)和一列predictions(元组输出的第二部分).但我不知道如何使用Python在Spark中写入CSV文件.

如何使用上述输出创建CSV文件?

Dan*_*bos 33

只是map在RDD(的线labelsAndPredictions)为字符串(该CSV的线),然后使用rdd.saveAsTextFile().

def toCSVLine(data):
  return ','.join(str(d) for d in data)

lines = labelsAndPredictions.map(toCSVLine)
lines.saveAsTextFile('hdfs://my-node:9000/tmp/labels-and-predictions.csv')
Run Code Online (Sandbox Code Playgroud)

  • 这取决于工具.Hadoop工具将读取所有`part-xxx`文件.当你使用`sc.textFile`时,Spark也会读取它.对于传统工具,您可能需要先将数据合并到一个文件中.如果输出小到足以通过传统工具处理,则没有理由通过Spark保存它.只需"收集"RDD并将数据写入没有Spark的本地文件. (7认同)
  • 我尝试使用它,但是当我执行它时会创建一个名为'labels-and-predictions.csv'的目录,并且在该目录中有两个文件--_SUCCESS和part-00000 (4认同)

Ins*_*ico 20

我知道这是一个老帖子.但是为了帮助搜索相同内容的人,以下是我如何在PySpark 1.6.2中将单列RDD写入单个CSV文件

RDD:

>>> rdd.take(5)
[(73342, u'cells'), (62861, u'cell'), (61714, u'studies'), (61377, u'aim'), (60168, u'clinical')]
Run Code Online (Sandbox Code Playgroud)

现在的代码:

# First I convert the RDD to dataframe
from pyspark import SparkContext
df = sqlContext.createDataFrame(rdd, ['count', 'word'])
Run Code Online (Sandbox Code Playgroud)

DF:

>>> df.show()
+-----+-----------+
|count|       word|
+-----+-----------+
|73342|      cells|
|62861|       cell|
|61714|    studies|
|61377|        aim|
|60168|   clinical|
|59275|          2|
|59221|          1|
|58274|       data|
|58087|development|
|56579|     cancer|
|50243|    disease|
|49817|   provided|
|49216|   specific|
|48857|     health|
|48536|      study|
|47827|    project|
|45573|description|
|45455|  applicant|
|44739|    program|
|44522|   patients|
+-----+-----------+
only showing top 20 rows
Run Code Online (Sandbox Code Playgroud)

现在写入CSV

# Write CSV (I have HDFS storage)
df.coalesce(1).write.format('com.databricks.spark.csv').options(header='true').save('file:///home/username/csv_out')
Run Code Online (Sandbox Code Playgroud)

PS:我只是一个初学者,从Stackoverflow中的帖子中学习.所以我不知道这是不是最好的方法.但它对我有用,我希望它会帮助别人!


Gal*_*ong 11

用逗号加入是不好的,因为如果字段包含逗号,则它们将不会被正确引用,例如','.join(['a', 'b', '1,2,3', 'c']),a,b,1,2,3,c在您需要时提供给您a,b,"1,2,3",c.相反,您应该使用Python的csv模块将RDD中的每个列表转换为格式正确的csv字符串:

# python 3
import csv, io

def list_to_csv_str(x):
    """Given a list of strings, returns a properly-csv-formatted string."""
    output = io.StringIO("")
    csv.writer(output).writerow(x)
    return output.getvalue().strip() # remove extra newline

# ... do stuff with your rdd ...
rdd = rdd.map(list_to_csv_str)
rdd.saveAsTextFile("output_directory")
Run Code Online (Sandbox Code Playgroud)

由于csv模块只写入文件对象,我们必须创建一个空的"文件",io.StringIO("")并告诉csv.writer将csv格式的字符串写入其中.然后,我们output.getvalue()用来获取我们刚刚写入"文件"的字符串.要使此代码适用于Python 2,只需将io替换为StringIO模块即可.

如果您正在使用Spark DataFrames API,您还可以查看具有csv格式的DataBricks保存功能.