【发布时间】: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