【问题标题】:filtering out spark dataframe using udf使用 udf 过滤掉 spark 数据帧
【发布时间】:2023-01-09 20:14:52
【问题描述】:

我有一个包含两列的 pyspark 数据框,name 和 source。 name 列中的所有值都是不同的。 Source 有多个字符串,用逗号 (,) 分隔。 我想过滤掉所有这些行,其中 source 列中的任何字符串包含整个 name 列中的任何值。

我正在使用以下 UDF:

def checkDependentKPI(df, name_list):
for row in df.collect():
    for src in row["source"].split(","):
        for name in name_list:
            if name in src:
                return row['name']
        return row['name']

我的最终目标是将所有此类行放在数据框的末尾。我该怎么做?

示例数据框:

+--------------------+--------------------+
|                name|              source|
+--------------------+--------------------+
|dev.................|prod, sum, diff.....|
|prod................|dev, diff, avg......|
|stage...............|mean, mode..........|
|balance.............|median, mean........|
|target..............|avg, diff, sum......|
+--------------------+--------------------+

【问题讨论】:

    标签: python pyspark


    【解决方案1】:

    使用split() 函数将您的 kpi_source_table 放入一个数组中,然后使用array_contains() 函数来过滤游览结果。

    from pyspark.sql import functions as psf 
    
    test_data = [('dev','prod,sum,diff')
      , ('prod','dev,diff,avg')
      , ('stage','mean,mode')
      , ('balance','median,mean')
      , ('target','avg,diff,sum')]
      
    df = spark.createDataFrame(test_data, ['kpi_name','kpi_source_table'])
    
    df = df.withColumn('kpi_source_table', psf.split('kpi_source_table', ','))
    df = df.withColumn('contains_avg', psf.array_contains(psf.col('kpi_source_table'), 'avg'))
    
    display(df.filter('contains_avg=True'))
    

    【讨论】:

    • 函数就在那里,我必须过滤两个数据帧,作为符合此条件的数据框和不符合此条件的数据框。 name_list 是名称列中所有值的列表。
    • 出于性能原因,您应该使用内置的 spark 函数而不是逐行操作来解决此问题
    • 我认为您不清楚问题,查看我的数据框,不同的名称是 dev、prod、stage、balance 和 target。我想将所有这些列放在数据框的末尾,其中源中的任何元素都包含名称列中的任何值,即 dev、prod、stage、balance 和 target。因此第 0 行和第 1 行应该在数据帧的末尾移动。
    猜你喜欢
    • 1970-01-01
    • 2017-09-16
    • 1970-01-01
    • 2018-09-17
    • 1970-01-01
    • 2020-09-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多