如何处理来自 blob 存储且数据块中路径较长的多个文件?

Chr*_*Wue 5 azure azure-blob-storage azure-databricks

我已启用 API 管理服务的日志记录,并且日志存储在存储帐户中。现在,我尝试在 Azure Databricks 工作区中处理它们,但在访问这些文件时遇到困难。

问题似乎是自动生成的虚拟文件夹结构如下所示:

/insights-logs-gatewaylogs/resourceId=/SUBSCRIPTIONS/<subscription>/RESOURCEGROUPS/<resource group>/PROVIDERS/MICROSOFT.APIMANAGEMENT/SERVICE/<api service>/y=*/m=*/d=*/h=*/m=00/PT1H.json
Run Code Online (Sandbox Code Playgroud)

我已将insights-logs-gatewaylogs容器安装在下面/mnt/diags,并dbutils.fs.ls('/mnt/diags')正确列出了该resourceId=文件夹,但未dbutils.fs.ls('/mnt/diags/resourceId=')找到声明文件

如果我沿着虚拟文件夹结构创建空标记 blob,我可以列出每个后续级别,但该策略显然会失败,因为路径的最后部分是按年/月/日/小时动态组织的。

例如一个

spark.read.format('json').load("dbfs:/mnt/diags/logs/resourceId=/SUBSCRIPTIONS/<subscription>/RESOURCEGROUPS/<resource group>/PROVIDERS/MICROSOFT.APIMANAGEMENT/SERVICE/<api service>/y=*/m=*/d=*/h=*/m=00/PT1H.json")
Run Code Online (Sandbox Code Playgroud)

此错误的产量:

java.io.FileNotFoundException: File/resourceId=/SUBSCRIPTIONS/<subscription>/RESOURCEGROUPS/<resource group>/PROVIDERS/MICROSOFT.APIMANAGEMENT/SERVICE/<api service>/y=2019 does not exist.
Run Code Online (Sandbox Code Playgroud)

很明显,通配符已经找到了第一年文件夹,但拒绝进一步向下。

我在 Azure 数据工厂中设置了一个复制作业,该作业成功复制同一 Blob 存储帐户中的所有 json Blob 并删除前缀resourceId=/SUBSCRIPTIONS/<subscription>/RESOURCEGROUPS/<resource group>/PROVIDERS/MICROSOFT.APIMANAGEMENT/SERVICE/<api service>(因此根文件夹以年份组件开头),并且可以一路成功访问,而无需创建空标记斑点。

因此,问题似乎与长虚拟文件夹结构有关,该结构大部分为空。

是否有另一种方法可以在 databricks 中处理此类文件夹结构?

更新:我也尝试在安装时提供路径作为安装的一部分source,但这也没有帮助

Chr*_*Wue 2

我想我可能已经找到了这个问题的根本原因。应该早点尝试过,但我提供了现有 blob 的确切路径,如下所示:

spark.read.format('json').load("dbfs:/mnt/diags/logs/resourceId=/SUBSCRIPTIONS/<subscription>/RESOURCEGROUPS/<resource group>/PROVIDERS/MICROSOFT.APIMANAGEMENT/SERVICE/<api service>/y=2019/m=08/d=20/h=06/m=00/PT1H.json")
Run Code Online (Sandbox Code Playgroud)

我得到了一个更有意义的错误:

shaded.databricks.org.apache.hadoop.fs.azure.AzureException:com.microsoft.azure.storage.StorageException:Blob 类型不正确,请使用正确的 Blob 类型访问服务器上的 Blob。预期为 BLOCK_BLOB,实际为 APPEND_BLOB。

事实证明,开箱即用的日志记录会创建附加 blob(并且似乎没有办法更改它),并且从这张票证的外观来看,对附加 blob 的支持仍然是 WIP:https://issues.apache .org/jira/browse/HADOOP-13475

FileNotFoundException可能是一个转移注意力的事情,可能是由于在尝试扩展通配符并查找不受支持的 blob 类型时被吞掉的内部异常引起的。

更新

终于找到了合理的解决方法。我azure-storage在工作区中安装了 Python 包(如果您在家使用 Scala,那么它已经安装了),并自行加载了 blob。下面的大多数代码都是为了添加通配支持,如果您愿意只匹配路径前缀,则不需要它:

%python

import re
import json
from azure.storage.blob import AppendBlobService


abs = AppendBlobService(account_name='<account>', account_key="<access_key>")

base_path = 'resourceId=/SUBSCRIPTIONS/<subscription>/RESOURCEGROUPS/<resource group>/PROVIDERS/MICROSOFT.APIMANAGEMENT/SERVICE/<api service>'
pattern = base_path + '/*/*/*/*/m=00/*.json'
filter = glob2re(pattern)

spark.sparkContext \
     .parallelize([blob.name for blob in abs.list_blobs('insights-logs-gatewaylogs', prefix=base_path) if re.match(filter, blob.name)]) \
     .map(lambda blob_name: abs.get_blob_to_bytes('insights-logs-gatewaylogs', blob_name).content.decode('utf-8').splitlines()) \
     .flatMap(lambda lines: [json.loads(l) for l in lines]) \
     .collect()
Run Code Online (Sandbox Code Playgroud)

glob2re由/sf/answers/2087468701/提供:

def glob2re(pat):
    """Translate a shell PATTERN to a regular expression.

    There is no way to quote meta-characters.
    """

    i, n = 0, len(pat)
    res = ''
    while i < n:
        c = pat[i]
        i = i+1
        if c == '*':
            #res = res + '.*'
            res = res + '[^/]*'
        elif c == '?':
            #res = res + '.'
            res = res + '[^/]'
        elif c == '[':
            j = i
            if j < n and pat[j] == '!':
                j = j+1
            if j < n and pat[j] == ']':
                j = j+1
            while j < n and pat[j] != ']':
                j = j+1
            if j >= n:
                res = res + '\\['
            else:
                stuff = pat[i:j].replace('\\','\\\\')
                i = j+1
                if stuff[0] == '!':
                    stuff = '^' + stuff[1:]
                elif stuff[0] == '^':
                    stuff = '\\' + stuff
                res = '%s[%s]' % (res, stuff)
        else:
            res = res + re.escape(c)
    return res + '\Z(?ms)'
Run Code Online (Sandbox Code Playgroud)

不太漂亮,但避免了数据的复制,并且可以封装在一个小实用程序类中。