【问题标题】:Multiple Filtering in PySparkPySpark 中的多重过滤
【发布时间】:2018-03-17 03:58:12
【问题描述】:

我已经将一个数据集导入到 Juputer notebook/PySpark 中通过 EMR 进行处理,例如:

data sample

我想在使用过滤器功能之前清理数据。这包括:

  1. 删除空白行或“0”或不适用的成本或日期。我认为过滤器类似于:.filter(lambda (a,b,c,d): b = ?, c % 1 == c, d = ?)。我不确定如何过滤水果并储存在这里。
  2. 删除不正确的值,例如“3”不是水果名称。这对于数字来说很容易(只是数字 % 1 == 数字),但我不确定它会如何过滤掉这些单词。
  3. 删除统计异常值的行,即与平均值相差 3 个标准差的行。所以这里的单元格 C4 显然需要删除,但我不确定如何将此逻辑合并到过滤器中。

我需要一次执行一个过滤器,还是有办法一次性过滤数据集(以 lambda 表示法)?

或者,编写一个 Spark SQL 查询是否更容易,而不是在“where”子句中有许多过滤器(但是上面的 #3 仍然难以用 SQL 编写)。

【问题讨论】:

标签: python pyspark elastic-map-reduce


【解决方案1】:

如果您在文档中阅读http://spark.apache.org/docs/2.1.0/api/python/pyspark.sql.html#pyspark.sql.DataFrame.filter,会写到

where() 是 filter() 的别名。

因此,您也可以安全地使用 'filter' 而不是 'where' 来处理多个条件。

编辑: 如果你想过滤很多列的很多条件,我更喜欢这种方法。

from dateutil.parser import parse
import pyspark.sql.functions as F

def is_date(string):
    try: 
       parse(string)
       return True
    except ValueError:
       return False
def check_date(d):
    if is_date(d):
        return d
    else:
        return None

date_udf = F.udf(check_date,StrinType())

def check_fruit(name):
    fruits_list #create a list of fruits(can easily find it on google)
                #difficult filtering words otherwise
                #try checking from what you want, rest will be filtered
    if name in fruits_list:
        return name
    else:
        return None

fruit_udf = F.udf(check_fruit,StringType())

def check_cost(value):
    mean, std #calculcated beforehand
    threshold_upper = mean + (3*std)
    threhold_lower = mean - (3*std)

    if value > threhold_lower and value < threshold_upper:
        return value
    else:
        return None
cost_udf = F.udf(check_cost,StringType())        

#Similarly create store_udf

df = df.select([date_udf(F.col('date')).alias('date'),\
            fruit_udf(F.col('fruit')).alias('fruit'),\
            cost_udf(F.col('cost')).alias('cost'),\
            store_udf(F.col('store')).alias('store')]).dropna()

这将导致所有列一起工作。

【讨论】:

  • 感谢您的链接。问题主要是如何在 #1/2/3 上面过滤(或使用 where),因为我不确定如何使用多个 where/filter 子句清理数据。
  • 这里的主要问题是我有一个需要在使用前清理的大数据集。如果我推断架构,它将导入不正确的值,但随后我必须编写一个复杂的 SQL 查询来删除这些实例(这是不可扩展的,因为有人可以将其他数据添加到具有更多错误的同一个文件中)。我尝试使用 .withColumn 语句将值转换为“类型”,但它也拒绝了这一点,因为数据丢失/不干净(例如,数字字段中的 customer_count 不会正确区分大小写)。因此,我被困在如何使用 Spark 中的可用功能来解决这个问题。
  • 是的,推断架构是不合适的。由于您在要过滤的许多列上有许多条件,因此我建议对所有列使用 UDF。检查答案中的编辑。
  • 更改架构(通过强制转换或其他方式)然后使用 df2 = df.filter(第 1 列的条件 | 第 2 列的条件 ...)也可以吗?
  • 当您的列包含不干净的值时更改架构将导致错误。例如,在列成本中,如果您有字符,您将无法将其转换为 FloatType。最好将所有值清除为 StringType,然后转换/转换为所需的架构。
猜你喜欢
  • 2021-05-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-11-29
  • 1970-01-01
  • 2022-01-13
  • 1970-01-01
  • 2019-06-22
相关资源
最近更新 更多