【问题标题】:Efficient data retrieval process between Azure Blob storage and Azure databricksAzure Blob 存储和 Azure Databricks 之间的高效数据检索过程
【发布时间】: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


    【解决方案1】:

    我建议不要列出文件,然后单独阅读它们,而是查看Azure Databricks Autoloader。它可能会使用通知来查找哪些新文件已上传到 Blob 存储,而不是列出文件。

    它也可以同时处理多个文件,而不是一个一个地读取它们并进行联合。

    如果你不需要连续处理数据,那么你可以使用.trigger(once=True)来模拟数据的批量加载。

    【讨论】:

    • 跟进您的回答。关于安全性,我更愿意包含任何类似秘密范围的东西。我看不出有任何可能将它包含在自动加载器中还是忽略了它?
    • 你能扩展你想要实现的目标吗?你想从秘密范围中包括什么?
    【解决方案2】:

    有多种方法可以做到这一点,但我会这样做:

    每当在您的 Azure 存储帐户中创建新的 Blob 时,使用 Azure Functions 触发您的 Python 代码。这将删除代码的轮询部分,并在您的存储帐户上有文件可用时立即将数据发送到数据块

    例如,对于近乎实时的报告,您可以使用 Azure 流分析并在 Event Hub 上运行查询并输出到 Power Bi。

    【讨论】:

    • 感谢您的输入!我对近乎实时的结构也有类似的想法。问题是我不完全确定使用 Databricks 作为持久存储在成本方面是否是一个好主意。对于 NRT 结构,在 Azure 函数记录文件后立即将其发送到流分析很可能是一个好主意。但是对于这些大量的数据,我不确定。不仅要看数据大小,因为我们谈论的是 3kb 大小的 avro 文件,这当然不是那么大,但我认为移动数据的交易成本会很高。
    猜你喜欢
    • 1970-01-01
    • 2020-11-15
    • 2019-06-09
    • 2019-03-08
    • 2019-08-06
    • 2021-07-11
    • 2019-07-01
    • 2019-07-28
    • 1970-01-01
    相关资源
    最近更新 更多