【问题标题】:Possible to take multiple input files and not create one RDD in pyspark?可以获取多个输入文件而不是在 pyspark 中创建一个 RDD?
【发布时间】:2017-10-12 22:29:54
【问题描述】:

在 Hadoop 中,我可以将应用程序指向一个路径,然后映射器将单独处理文件。我必须以这种方式处理它,因为我需要解析文件名和路径以匹配我直接在映射器中加载的其他文件。

在 pyspark 中,将路径传递给 SparkContext 的 textFile 会创建一个 RDD。有没有办法在 Spark / pyspark 中复制相同的 Hadoop 行为?

【问题讨论】:

    标签: hadoop pyspark


    【解决方案1】:

    我希望这能解决您的一些困惑: sparkContext.wholeTextFiles(path) 返回 pairRDD(有用链接:https://www.safaribooksonline.com/library/view/learning-spark/9781449359034/ch04.html

    简而言之,pairRDD 更像是一张地图(即有键、有值)

    rdd = sparkContext.wholeTextFiles(path)
    
    def func_work_on_individual_files(x):
       # x is a tuple which will receive both (key, value) for the pairRDD Row Elements passed. key -> file path, value -> content of a file with line seperated by '/n' (as you mentioned). To access key use x[0], to access value use x[1]. 
       # your logic to do something useful with file data, 
       # to get separate lines you can use: x[1].split('\n')
       # end function by return the values you want to return out of a file's data. 
    
       # I am simply returning the whole content of file 
       return x[1] 
    
    
    #loop over each of the file in the pairRdd created above
    file_contents = rdd.map(func_work_on_individual_files)
    
    #this will create just one partition out of all elements in list (as you mentioned)
    consolidated_contents = file_contents.repartition(1)
    
    #Save final output - this will create just one path like Hadoop
    consolidated_contents.saveAsTextFile(path)
    

    【讨论】:

    • 谢谢!但是在那个过程中它不再是一个 rdd 而是一个列表?由于尺寸,我需要将其保留为 rdd / 不想在运输过程中丢失它。我可能会改变输入来解决这个问题。
    • 标记为答案,但处理方式略有不同。我需要通过索引将数据与其他内容配对,这变得太有问题了,因为我找不到防止数据成为非 RDD 对象的方法。我现在只是使用 wholeTextFiles 作为获取所有文件名的一种方式。不太喜欢这个解决方案,但它会起作用。感谢您的帮助!它确实有助于澄清一些事情
    • 直到你不使用一个动作,它仍然是一个rdd。在我的示例代码中,唯一的操作是“saveAsTextFile”。所有其他操作都是转换,即不创建列表或任何其他明确的数据结构。我仍然不清楚您的用例,但是如果您想用每条记录(数据的每一行)标记文件名,您可以从“func_work_on_individual_files”返回一个键值对,它将保留它一个pairRDD。或者您可以以多种方式对每条记录进行反规范化,例如:使用分隔符将文件名标记到每条记录(每行)并返回索引,记录来自 udf。
    【解决方案2】:

    Pyspark 为这个用例提供了一个函数:sparkContext.wholeTextFiles(path)。它将读取一个文本文件目录并生成一个键值对,其中 key 是每个文件的路径,value 是每个文件的内容。

    【讨论】:

    • 您将如何遍历文件?我尝试做这样的事情......“files.collect().foreach(lambda文件:”......但任何分配都失败了。文件是字符串文件名的rdd。似乎可以在Scala-spark中做到。
    • 另外,我想如果我设法对其进行迭代,那么输出将是单独的,需要不同的输出文件夹吗?在 Hadoop 中,它将是一个输出(按部分文件很好地分割,但可以合并到一个文件中)
    • 我一直试图弄清楚如何迭代它,但注意到这个方法实际上为输入返回了一个非常不同的结构。 textFile 将返回列表列表,但 wholeTextFiles 返回“content”,其中“\r”和“\n”作为分隔符。真的不明白为什么会有不同的处理方式。
    猜你喜欢
    • 2016-06-07
    • 1970-01-01
    • 2020-09-22
    • 2017-04-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-09-06
    • 1970-01-01
    相关资源
    最近更新 更多