【问题标题】:how to upload a parquet file into Azure ADLS 2 Blob如何将镶木地板文件上传到 Azure ADLS 2 Blob
【发布时间】:2020-09-20 15:09:47
【问题描述】:

您好,我想将 Parquet 文件上传到 ADLS gen 2 blob。我正在使用下面的代码行来创建 blob 并在其中上传 parquet 文件。

blob = BlobClient.from_connection_string(conn_str="Connection String", container_name="parquet", blob_name=outdir)
df.to_parquet('logs.parquet',compression='GZIP') #df is dataframe
with open("./logs.parquet", "rb") as data:
blob.upload_blob(data)
os.remove("logs.parquet")

我没有遇到任何错误,文件也写在 blob 中。但是,我不认为我做对了,因为 ADX/kusto 查询无法理解文件并且那里没有数据可见。

以下是我在 Azure 数据资源管理器中执行的步骤,以从 ADLS gen 2 中上传的 parquet 文件中获取记录。

已创建外部表:

.create external table LogDataParquet(AppId_s:string,UserId_g:string,Email_s:string,RoleName_s:string,Operation_s:string,EntityId_s:string,EntityType_s:string,EntityName_s:string,TargetTitle_s:string,TimeGenerated:datetime) 
kind=blob
dataformat=parquet
( 
   h@'https://streamoutalds2.blob.core.windows.net/stream-api-raw-testing;secret'
)
with 
(
   folder = "ExternalTables"   
)

外部表列映射:

.create external table LogDataParquet parquet mapping "LogDataMapparquet" '[{ "column" : "AppId_s", "path" : "$.AppId_s"},{ "column" : "UserId_g", "path" : "$"},{ "column" : "Email_s", "path" : "$.Email_s"},{ "column" : "RoleName_s", "path" : "$.RoleName_s"},{ "column" : "Operation_s", "path" : "$.Operation_s"},{ "column" : "EntityId_s", "path" : "$.EntityId_s"}]'

外部表不提供记录

external_table('LogDataParquet')

没有记录

external_table('LogDataParquet') | count 

1 条记录 - 计数 0

我使用流分析使用了类似的场景,我接收传入的流并将其以镶木地板格式保存到 ADLS。在这种情况下,ADX 中的外部表可以很好地获取记录。我觉得我在用 blob 编写镶木地板文件的方式上犯了错误 - (使用 open("./logs.parquet", "rb") 作为数据:)

【问题讨论】:

  • 您能否提供更多有关您在 Kusto 查询中看到的错误的详细信息?
  • 我没有看到 kusto 有任何错误,但它没有获取任何记录。只是空行。
  • 如果您看到空行(并且空行数与 parquet 文件中的输入行数匹配) - 这可能表明您的摄取映射不正确。该问题未指定数据在上传到 blob 后如何进入 ADX。您是使用 EventGrid 还是使用 ADX python 库来获取数据?请澄清一下。
  • 空行,计数也为 0。我正在使用 ADX python 库来摄取数据(从 azure.storage.blob 导入 ContainerClient、BlobClient)。我已经用更多细节编辑了这个问题。

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


【解决方案1】:

根据日志,外部表定义如下:

.create external table LogDataParquet(AppId_s:string,UserId_g:string,Email_s:string,RoleName_s:string,Operation_s:string,EntityId_s:string,EntityType_s:string,EntityName_s:string,TargetTitle_s:string,TimeGenerated:datetime) 
kind=blob
partition by 
   AppId_s,
   bin(TimeGenerated,1d)
dataformat=parquet
( 
   '******'
)
with 
(
   folder = "ExternalTables"   
)

PARTITION BY 子句告诉 ADX 预期的文件夹布局是:

<AppId_s>/<TimeGenerated, formatted as 'yyyy/MM/dd'>

例如:

https://streamoutalds2.blob.core.windows.net/stream-api-raw-testing;secret/SuperApp/2020/01/31

您可以在此部分中找到有关 ADX 如何在查询期间定位外部存储上的文件的更多信息:https://docs.microsoft.com/en-us/azure/data-explorer/kusto/management/external-tables-azurestorage-azuredatalake#artifact-filtering-logic

要根据文件夹布局修复外部表定义,请使用.alter 命令:

.alter external table LogDataParquet(AppId_s:string,UserId_g:string,Email_s:string,RoleName_s:string,Operation_s:string,EntityId_s:string,EntityType_s:string,EntityName_s:string,TargetTitle_s:string,TimeGenerated:datetime) 
kind=blob
dataformat=parquet
( 
  h@'https://streamoutalds2.blob.core.windows.net/stream-api-raw-testing;secret'
)
with 
(
   folder = "ExternalTables"   
)

顺便说一句,如果映射是幼稚的(例如,映射的列名与数据源列名匹配),那么 Parquet 格式就不需要它。

【讨论】:

  • 是的,表定义期望具有“App_id 和日期”的数据分区,并且在 ADLS blob 中也维护相同。我没有提到分区,因为我在不同的场景中使用过它,并且认为问题不存在。此外,即使我从外部表中删除所有分区定义并从 ADLS blob 中提供镶木地板文件的完整路径,表仍然不会获取任何记录。我没有遇到由 Azure Stream Analytics 编写的一组镶木地板文件的任何问题,但仅限于 Python SDK。
  • 你可以在某处上传示例 Parquet 文件吗?
  • @ashwiniprakash 我无权访问数据。如果您只是在某些公共服务上上传示例 Parquet 文件,那将真的很有帮助。谢谢!
  • 请在我的 git 存储库 'github.com/ashwini-git/datastore.git' 中找到 parquet 文件。谢谢!
  • 我用 pyarrow 代替了 fastparquet 和它的工作原理。
猜你喜欢
  • 2020-07-27
  • 1970-01-01
  • 1970-01-01
  • 2021-10-11
  • 2022-01-15
  • 1970-01-01
  • 2023-03-27
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多