【问题标题】:Is there a way to load multiple text files into a single dataframe using Databricks?有没有办法使用 Databricks 将多个文本文件加载到单个数据框中?
【发布时间】:2019-10-02 15:56:28
【问题描述】:

我正在尝试测试一些想法,以递归循环遍历文件夹和子文件夹中的所有文件,并将所有内容加载到单个数据框中。我有 12 种不同的文件,不同之处在于文件命名约定。所以,我有以“ABC”开头的文件名、以“CN”开头的文件名、以“CZ”开头的文件名等等。我尝试了以下 3 个想法。

import pyspark  
import os.path
from pyspark.sql import SQLContext
from pyspark.sql.functions import input_file_name

df = sqlContext.read.format("com.databricks.spark.text").option("header", "false").load("dbfs/mnt/rawdata/2019/06/28/Parent/ABC*.gz")
df.withColumn('input', input_file_name())
print(dfCW)

df = sc.textFile('/mnt/rawdata/2019/06/28/Parent/ABC*.gz')
print(df)

df = sc.sequenceFile('dbfs/mnt/rawdata/2019/06/28/Parent/ABC*.gz/').toDF()
df.withColumn('input', input_file_name())
print(dfCW)

这可以通过 PySpark 或 PySpark SQL 完成。我只需要将数据湖中的所有内容加载到数据框中,这样我就可以将数据框推送到 Azure SQL Server 中。我在 Azure Databricks 中进行所有编码。如果这是普通的 Python,我可以很容易地做到这一点。我只是不太了解 PySpark,无法使其正常工作。

为了说明这一点,我有 3 个如下所示的压缩文件(ABC0006.gz、ABC00015.gz 和 ABC0022.gz):

ABC0006.gz
0x0000fa00|ABC|T3|1995
0x00102c55|ABC|K2|2017
0x00024600|ABC|V0|1993

ABC00015.gz
0x00102c54|ABC|G1|2016
0x00102cac|ABC|S4|2017
0x00038600|ABC|F6|2003

ABC0022.gz
0x00102c57|ABC|J0|2017
0x0000fa00|ABC|J6|1994
0x00102cec|ABC|V2|2017

我想将所有内容合并到一个看起来像这样的 datdframe(.gz 是文件名;每个文件都有完全相同的标题):

0x0000fa00|ABC|T3|1995
0x00102c55|ABC|K2|2017
0x00024600|ABC|V0|1993
0x00102c54|ABC|G1|2016
0x00102cac|ABC|S4|2017
0x00038600|ABC|F6|2003
0x00102c57|ABC|J0|2017
0x0000fa00|ABC|J6|1994
0x00102cec|ABC|V2|2017

我有 1000 多个这样的文件要处理。幸运的是,只有 12 种不同类型的文件,因此有 12 种名称……以“ABC”、“CN”、“CZ”等开头。感谢您查看此处。

根据您的 cmets,Abraham,看来我的代码应该是这样的,对吧...

file_list=[]
path = 'dbfs/rawdata/2019/06/28/Parent/'
files  = dbutils.fs.ls(path)
for file in files:
    if(file.name.startswith('ABC')):
       file_list.append(file.name)
df = spark.read.load(path=file_list)

这是正确的,还是不正确的?请指教。我认为我们很接近,但这仍然对我不起作用,否则我不会在这里重新发布。谢谢!!

【问题讨论】:

    标签: dataframe pyspark apache-spark-sql pyspark-sql azure-databricks


    【解决方案1】:

    PySpark 支持使用 load 函数加载文件列表。我相信这就是你要找的东西

    file_list=[]
    path = 'dbfs/mnt/rawdata/2019/06/28/Parent/'
    files  = dbutils.fs.ls(path)
    for file in files:
        if(file.name.startswith('ABC')):
           file_list.append(file.name)
    df = spark.read.load(path=file_list)
    

    如果文件是 CSV 并且有标题,请使用以下命令

    df = spark.read.load(path=file_list,format="csv", sep=",", inferSchema="true", header="true")
    

    更多示例代码请参考https://spark.apache.org/docs/latest/sql-data-sources-load-save-functions.html

    【讨论】:

    • 感谢亚伯拉罕!这看起来很有希望!!我将路径更改为我的实际路径,运行代码,现在我收到此错误消息:NameError: name 'dbfsutils' is not defined。我导入了 DBUtils 库,看起来它导入得很好。是否需要任何其他依赖项才能完成这项工作?
    • @asher 感谢您的指出。那是一个错字。执行代码不需要其他依赖项 dbutils 和 spark session。
    • 好的,代码在 Python 环境下看起来不错;我在 Databricks 上的 PySpark 中运行它。我取出了“import DBUtils”这一行,并从集群中删除了该库。我刚刚重新运行了代码,得到了相同的结果: NameError: name 'dbfsutils' is not defined 此行发生错误:files = dbfsutils.fs.ls(path)
    • @asher 只需将 dbfsutils 更改为 dbutils。我在 Databricks 中尝试过,它给出了预期的结果
    【解决方案2】:

    我终于,终于,终于搞定了。

    val myDFCsv = spark.read.format("csv")
       .option("sep","|")
       .option("inferSchema","true")
       .option("header","false")
       .load("mnt/rawdata/2019/01/01/client/ABC*.gz")
    
    myDFCsv.show()
    myDFCsv.count()
    

    显然所有压缩文件和推断模式任务都是自动处理的。因此,代码超级、超级轻量级​​,而且速度也非常快。

    【讨论】:

      猜你喜欢
      • 2022-10-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-03-20
      • 1970-01-01
      相关资源
      最近更新 更多