【发布时间】:2021-03-22 20:46:12
【问题描述】:
如果所有值都是 ASCII,我需要检查 pyspark 数据帧,我使用以下方法:
def is_ascii(s):
if s:
return all(ord(c) < 128 for c in s)
else:
return None
is_ascii_udf = udf(lambda l: is_ascii(l), BooleanType() )
df_result = df.select( *map(lambda col: is_ascii_udf(df[col]).alias(col), df.columns ) )
我正在尝试将其与具有 50MM 行和 9000 列的新数据一起使用,但出现此错误:
ExecutorLostFailure (executor 30 exited caused by one of the running tasks) Reason: Remote RPC client disassociated. Likely due to containers exceeding thresholds, or network issues.
看来内存已满,无法获得更大的集群,所以我想这样做
import pyspark.sql.functions as F
import pandas as pd
from pyspark.sql.types import *
df = spark.read.parquet( path)
for i in df.columns:
df = spark.read.parquet( path)
df_result = df.select( *map(lambda col: is_ascii_udf(df[col]).alias(col), [i] ) )
n = df_result.filter( ~F.col(i) ).count()
if n>0:
print(i,n)
但是我得到同样的错误,为什么我每次读取数据帧并且只对一列执行 udf 时仍然得到同样的错误
集群有 50 GB 内存,6 个核心,最多 8 个工作人员
我认为错误在于函数,或者我如何使用它
问候
【问题讨论】:
标签: python-3.x pyspark ascii user-defined-functions pyspark-dataframes