【问题标题】:Bulk upload to Azure Data Lake Gen 2 with REST APIs使用 REST API 批量上传到 Azure Data Lake Gen 2
【发布时间】:2020-01-19 09:32:57
【问题描述】:

在另一个 related question 中,我曾询问如何将文件从本地上传到 Microsoft Azure Data Lake Gen 2,通过 REST API 提供了答案。为了完整起见,建议的代码可以在下面找到。

由于对于大量相对较小的文件(0.05 MB),这种顺序上传文件已被证明相对较慢,我想问一下是否有可能一次对所有文件进行批量上传,假设所有文件的路径都是事先知道的?

使用 REST API 将单个文件上传到 ADLS Gen 2 的代码:

import requests
import json

def auth(tenant_id, client_id, client_secret):
    print('auth')
    auth_headers = {
        "Content-Type": "application/x-www-form-urlencoded"
    }
    auth_body = {
        "client_id": client_id,
        "client_secret": client_secret,
        "scope" : "https://storage.azure.com/.default",
        "grant_type" : "client_credentials"
    }
    resp = requests.post(f"https://login.microsoftonline.com/{tenant_id}/oauth2/v2.0/token", headers=auth_headers, data=auth_body)
    return (resp.status_code, json.loads(resp.text))

def mkfs(account_name, fs_name, access_token):
    print('mkfs')
    fs_headers = {
        "Authorization": f"Bearer {access_token}"
    }
    resp = requests.put(f"https://{account_name}.dfs.core.windows.net/{fs_name}?resource=filesystem", headers=fs_headers)
    return (resp.status_code, resp.text)

def mkdir(account_name, fs_name, dir_name, access_token):
    print('mkdir')
    dir_headers = {
        "Authorization": f"Bearer {access_token}"
    }
    resp = requests.put(f"https://{account_name}.dfs.core.windows.net/{fs_name}/{dir_name}?resource=directory", headers=dir_headers)
    return (resp.status_code, resp.text)

def touch_file(account_name, fs_name, dir_name, file_name, access_token):
    print('touch_file')
    touch_file_headers = {
        "Authorization": f"Bearer {access_token}"
    }
    resp = requests.put(f"https://{account_name}.dfs.core.windows.net/{fs_name}/{dir_name}/{file_name}?resource=file", headers=touch_file_headers)
    return (resp.status_code, resp.text)

def append_file(account_name, fs_name, path, content, position, access_token):
    print('append_file')
    append_file_headers = {
        "Authorization": f"Bearer {access_token}",
        "Content-Type": "text/plain",
        "Content-Length": f"{len(content)}"
    }
    resp = requests.patch(f"https://{account_name}.dfs.core.windows.net/{fs_name}/{path}?action=append&position={position}", headers=append_file_headers, data=content)
    return (resp.status_code, resp.text)

def flush_file(account_name, fs_name, path, position, access_token):
    print('flush_file')
    flush_file_headers = {
        "Authorization": f"Bearer {access_token}"
    }
    resp = requests.patch(f"https://{account_name}.dfs.core.windows.net/{fs_name}/{path}?action=flush&position={position}", headers=flush_file_headers)
    return (resp.status_code, resp.text)

def mkfile(account_name, fs_name, dir_name, file_name, local_file_name, access_token):
    print('mkfile')
    status_code, result = touch_file(account_name, fs_name, dir_name, file_name, access_token)
    if status_code == 201:
        with open(local_file_name, 'rb') as local_file:
            path = f"{dir_name}/{file_name}"
            content = local_file.read()
            position = 0
            append_file(account_name, fs_name, path, content, position, access_token)
            position = len(content)
            flush_file(account_name, fs_name, path, position, access_token)
    else:
        print(result)


if __name__ == '__main__':
    tenant_id = '<your tenant id>'
    client_id = '<your client id>'
    client_secret = '<your client secret>'

    account_name = '<your adls account name>'
    fs_name = '<your filesystem name>'
    dir_name = '<your directory name>'
    file_name = '<your file name>'
    local_file_name = '<your local file name>'

    # Acquire an Access token
    auth_status_code, auth_result = auth(tenant_id, client_id, client_secret)
    access_token = auth_status_code == 200 and auth_result['access_token'] or ''
    print(access_token)

    # Create a filesystem
    mkfs_status_code, mkfs_result = mkfs(account_name, fs_name, access_token)
    print(mkfs_status_code, mkfs_result)

    # Create a directory
    mkdir_status_code, mkdir_result = mkdir(account_name, fs_name, dir_name, access_token)
    print(mkdir_status_code, mkdir_result)

    # Create a file from local file
    mkfile(account_name, fs_name, dir_name, file_name, local_file_name, access_token)

【问题讨论】:

  • 如果想批量上传很多文件到azure data lake gen 2,性能不错,可以使用python调用azcopy。请参阅我的更新答案。

标签: python azure azure-storage azure-data-lake


【解决方案1】:

到目前为止,将大量文件上传到 ADLS gen2 的最快方法是使用 AzCopy。您可以编写 python 代码来调用 AzCopy。

首先,按照link下载AzCopy.exe,下载后将文件解压,将azcopy.exe复制到一个文件夹(无需安装,是可执行文件),如F:\\azcopy\\v10\\azcopy.exe

然后从 azure 门户生成 sas 令牌,然后复制并保存 sas 令牌:

假设您已经为您的 adls gen2 帐户创建了文件系统,并且您不需要手动创建目录,它将由 azcopy 自动创建。

您需要注意的另一件事是,对于端点,您应该使用将 dfs 更改为 blob:例如将 https://youraccount.dfs.core.windows.net/ 更改为 https://youraccount.blob.core.windows.net/。

示例代码如下:

import subprocess

exepath = "F:\\azcopy\\v10\\azcopy.exe"
local_directory="F:\\temp\\1\\*"
sasToken="?sv=2018-03-28&ss=bfqt&srt=sco&sp=rwdlacup&se=2019-09-20T09:44:22Z&st=2019-09-20T01:44:22Zxxxxxxxx"

#note for the endpoint, you should change dfs to blob
endpoint="https://yygen2.blob.core.windows.net/w22/testfile5/"
myscript=exepath + " copy " + "\""+ local_directory + "\" " + "\""+endpoint+sasToken + "\"" + " --recursive"

print(myscript)
subprocess.call(myscript)

print("completed")

测试结果如下,本地目录下的所有文件/子文件夹都上传到ADLS gen2:

【讨论】:

  • 我自己实现了for循环,但是对于大量的小文件性能很差。但是我注意到 REST API 中没有批量上传功能。
  • @AlexGuevara,另外一种方式是你可以使用azcopy,它支持上传文件夹,可以提供更好的性能。您可以使用 python 代码调用 azcopy。并使用此信息更新了帖子。
猜你喜欢
  • 2023-03-04
  • 2020-11-22
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-11-10
  • 2019-09-23
  • 2020-01-24
  • 2020-09-27
相关资源
最近更新 更多