【问题标题】:Check ASCII pyspark dataframe检查 ASCII pyspark 数据帧
【发布时间】: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


    【解决方案1】:

    即使在一列上运行它也可能对您的集群来说太多了。不管怎样,有 Spark SQL 方法可以做你想做的事,在性能和内存方面应该更有效。

    下面的代码将给出每列中非 ascii 字符的布尔值或计数,并将结果收集到一个列表中。

    df.createOrReplaceTempView('df')
    
    is_not_ascii = [[col, spark.sql('select max(array_max(transform(split(%s, ""), x -> ascii(x))) >= 128) as is_ascii from df' % col).collect()[0][0]] for col in df.columns]
    # e.g. [['key', False], ['val', False]]
    
    count_not_ascii = [[col, spark.sql('select sum(cast(array_max(transform(split(%s, ""), x -> ascii(x))) >= 128 as int)) as is_ascii from df' % col).collect()[0][0]] for col in df.columns]
    # e.g. [['key', 0], ['val', 0]]
    

    【讨论】:

      猜你喜欢
      • 2021-08-18
      • 2020-07-10
      • 2022-01-08
      • 1970-01-01
      • 2021-11-12
      • 2017-03-16
      • 2021-04-10
      • 2018-10-23
      • 2020-06-10
      相关资源
      最近更新 更多