【发布时间】:2021-12-26 15:42:23
【问题描述】:
继续Managing huge zip files in dataBricks
Databricks 在 30 个文件后挂起。怎么办?
我已将巨大的 32Gb zip 拆分为 100 个独立的部分。我已经从文件中拆分了标题,因此可以像处理任何 CSV 文件一样处理它。我需要根据列过滤数据。文件位于 Azure Data Lake Storage Gen1 中,必须存储在那里。
尝试一次读取单个文件(或所有 100 个文件)在工作约 30 分钟后失败。 (请参阅上面的链接问题。)
我做了什么:
def lookup_csv(CR_nro, hlo_lista =[], output = my_output_dir ):
base_lib = 'adl://azuredatalakestore.net/<address>'
all_files = pd.DataFrame(dbutils.fs.ls(base_lib + f'CR{CR_nro}'), columns = ['full', 'name', 'size'])
done = pd.DataFrame(dbutils.fs.ls(output), columns = ['full', 'name', 'size'])
all_files = all_files[~all_files['name'].isin(tehdyt['name'].str.replace('/', ''))]
all_files = all_files[~all_files['name'].str.contains('header')]
my_scema = spark.read.csv(base_lib + f'CR{CR_nro}/header.csv', sep='\t', header=True, maxColumns = 1000000).schema
tmp_lst = ['CHROM', 'POS', 'ID', 'REF', 'ALT', 'QUAL', 'FILTER', 'INFO', 'FORMAT'] + [i for i in hlo_lista if i in my_scema.fieldNames()]
for my_file in all_files.iterrows():
print(my_file[1]['name'], time.ctime(time.time()))
data = spark.read.option('comment', '#').option('maxColumns', 1000000).schema(my_scema).csv(my_file[1]['full'], sep='\t').select(tmp_lst)
data.write.csv( output + my_file[1]['name'], header=True, sep='\t')
这行得通……有点。它的工作原理是大约 30 个文件,然后挂断
Py4JJavaError:调用 o70690.csv 时出错。 原因:org.apache.spark.SparkException:作业因阶段失败而中止:阶段 154.0 中的任务 0 失败 4 次,最近一次失败:阶段 154.0 中丢失任务 0.3(TID 1435、10.11.64.46、执行程序 7):com .microsoft.azure.datalake.store.ADLException:创建文件
CR03_pt29.vcf.gz/_started_1438828951154916601 时出错 操作 CREATE 因 HTTP401 失败:null 2 次尝试后最后遇到的异常抛出。 [HTTP401(null),HTTP401(null)]
我尝试添加一些删除并休眠:
data.unpersist()
data = []
time.sleep(5)
还有一些 try-exception 尝试。
for j in range(1,24):
for i in range(4):
try:
lookup_csv(j, hlo_lista =FN_list, output = blake +f'<my_output>/CR{j}/' )
except Exception as e:
print(i, j, e)
time.sleep(60)
这些都不走运。一旦失败,就会一直失败。
知道如何处理这个问题吗?我认为与 ADL 驱动器的连接会在一段时间后失败,但如果我将命令排队:
lookup_csv(<inputs>)
<next cell>
lookup_csv(<inputs>)
它可以正常工作,失败并且可以在下一个单元格中正常工作。我可以忍受这一点,但非常烦人的是基本循环无法在这种环境中工作。
【问题讨论】:
-
您是否使用启用了直通的集群?
-
@AlexOtt 是的,就是这样。
-
哦...docs.microsoft.com/en-us/azure/databricks/security/… 引用“您无法使用 Azure Active Directory 令牌生命周期策略延长 Azure Active Directory 直通令牌的生命周期。因此,如果您向集群发送命令,该命令需要超过一个小时,如果在 1 小时后访问 Azure Data Lake Storage 资源,它将失败。” /引用。好吧,如果我没看错的话,那就是“#¤”#¤。
-
是的。那是问题。每次调用笔记本单元格时都会刷新令牌,这就是为什么执行下一个时不会出错的原因。但是如果一个单元格花费的时间超过一小时,那么任务就会失败
-
@AlexOtt 谢谢...我想我只能循环遍历单元格。或运行包含多个单元格的工作簿。愚蠢,愚蠢的限制。
标签: apache-spark pyspark azure-databricks file-management