如何在Databricks笔记本中查看PySpark采集数据的前10条记录
问题背景
现有可正常运行的代码实现了从Azure存储账户读取JSON日志文件的功能,需要在Databricks笔记本中验证输出是否符合预期。由于文件数据量较大,无需拉取全量数据,仅需查看前10条记录即可完成校验。
原代码如下:
import re import json %pip install azure import azure from azure.storage.blob import AppendBlobService abs = AppendBlobService(account_name="azurestorage", account_key="mykey") base_path = "resourceId=/SUBSCRIPTIONS/5315MyId/RESOURCEGROUPS/AZURE-DEV/PROVIDERS/MICROSOFT.CONTAINERSERVICE/MANAGEDCLUSTERS/AZURE-DEV/y=2022/m=05/d=23/h=13/m=00/PT1H.json" pattern = base_path + "/*/*/*/*/m=00/*.json" filter = glob2re(pattern) df1 = ( spark.sparkContext.parallelize( [ blob.name for blob in abs.list_blobs("insights-logs-kube-audit", prefix=base_path) if re.match(filter, blob.name) ] ) .map( lambda blob_name: abs.get_blob_to_bytes("insights-logs-kube-audit", blob_name) .content.decode("utf-8") .splitlines() ) .flatMap(lambda lines: [json.loads(l) for l in lines]) .collect() )
实现方法
原代码末尾使用collect()会把所有解析完成的数据全量拉取到Driver节点,数据量大时极易触发Driver内存溢出,且全量拉取耗时很长,完全没必要用于抽样校验。直接替换为take(10)即可,该操作只会从集群中返回前10条解析完成的记录,执行效率高且资源占用极低。
另外注意两个代码细节:
- Databricks中
%pip安装命令需要放在单元格最开头,不能和import语句混放,建议把安装逻辑单独放在一个单元格执行 - 原代码调用了
glob2re方法但未导入,需要补充对应导入语句,否则会抛名称错误
修改后的可运行代码如下:
# 单独单元格先执行依赖安装 %pip install azure
# 后续单元格执行业务逻辑 import re import json from fnmatch import translate as glob2re from azure.storage.blob import AppendBlobService abs = AppendBlobService(account_name="azurestorage", account_key="mykey") base_path = "resourceId=/SUBSCRIPTIONS/5315MyId/RESOURCEGROUPS/AZURE-DEV/PROVIDERS/MICROSOFT.CONTAINERSERVICE/MANAGEDCLUSTERS/AZURE-DEV/y=2022/m=05/d=23/h=13/m=00/PT1H.json" pattern = base_path + "/*/*/*/*/m=00/*.json" file_filter = glob2re(pattern) # 仅取前10条数据 top10_records = ( spark.sparkContext.parallelize( [ blob.name for blob in abs.list_blobs("insights-logs-kube-audit", prefix=base_path) if re.match(file_filter, blob.name) ] ) .map( lambda blob_name: abs.get_blob_to_bytes("insights-logs-kube-audit", blob_name) .content.decode("utf-8") .splitlines() ) .flatMap(lambda lines: [json.loads(l) for l in lines]) .take(10) ) # Databricks中直接调用display即可格式化渲染输出,比print可读性更好 display(top10_records)
执行后就可以在笔记本输出区域直接看到格式化后的前10条JSON记录,快速校验字段、格式是否符合预期。
注:原代码中用filter做变量名会覆盖Python内置的filter函数,修改后重命名为file_filter避免潜在问题。
内容的提问来源于stack exchange,提问作者ZZZSharePoint
相关产品推荐
相关产品推荐

