【发布时间】:2021-07-31 21:04:53
【问题描述】:
我正在设计一个新的数据环境,目前正在开发我的概念证明。在这里,我使用以下架构: Azure 函数 --> Azure 事件中心 --> Azure Blob 存储 --> Azure 工厂 --> Azure 数据块 --> Azure SQL 服务器。
我目前正在努力解决的问题是如何优化“数据检索”以在 Azure Databricks 上提供我的 ETL 流程。
我正在处理通过前面的通道按分钟提交到 Azure blob 存储的事务性工厂数据。因此,我每天需要处理 86000 个文件。实际上,这是要处理的大量单独文件。目前,我使用以下代码来构建当前存在于 azure blob 存储中的文件名列表。接下来,我通过使用循环读取每个文件来检索它们。
我面临的问题是这个过程需要的时间。当然,我们在这里谈论的是需要读取的大量小文件。所以我不指望这个过程会在几分钟内完成。
我知道升级 databricks 集群可能会解决问题,但我不确定只有这样才能解决问题,查看在这种情况下我需要传输的文件数量。我正在运行 databricks 的以下代码。
# Define function to list content of mounted folder
def get_dir_content(ls_path):
dir_paths = ""
dir_paths = dbutils.fs.ls(ls_path)
subdir_paths = [get_dir_content(p.path) for p in dir_paths if p.isDir() and p.path != ls_path]
flat_subdir_paths = [p for subdir in subdir_paths for p in subdir]
return list(map(lambda p: p.path, dir_paths)) + flat_subdir_paths
filenames = []
paths = 0
mount_point = "PATH"
paths = get_dir_content(mount_point)
for p in paths:
# print(p)
filenames.append(p)
avroFile = pd.DataFrame(filenames)
avroFileList = avroFile[(avroFile[0].str.contains('.avro')) & (avroFile[0].str.contains('dbfs:/mnt/PATH'))]
avro_result = []
# avro_file = pd.DataFrame()
avro_complete = pd.DataFrame()
for i in avroFileList[0]:
avro_file = spark.read.format("avro").load(i)
avro_result.append(avro_file)
最后,我对所有这些文件进行联合,以创建它们的一个数据框。
# Schema definiëren op basis van
avro_df = avro_result[0]
# Union all dataframe
for i in avro_result:
avro_df = avro_df.union(i)
display(avro_df)
我想知道如何优化这个过程。按分钟输出的原因是,一旦我们有了分析报告架构(我们只需要每天的流程),我们计划稍后构建“近乎实时的洞察力”。
【问题讨论】:
标签: azure pyspark azure-blob-storage databricks azure-databricks