使用Pyspark/Pandas逐行统计多列相同值数量及对应取值
逐行多列频次统计实现
你可以根据使用的工具栈选择Pandas或PySpark实现对应逻辑:
规则说明
源数据固定包含A、B、C、D四个数值列,逐行计算后输出两个新列:
#_identical:当前行四个列值中,出现次数最高的频次value:对应最高频次的取值,仅一个最高频值时存该值,多个值频次并列最高时存所有值的列表
计算逻辑参考样例:
| A | B | C | D | #_identical | value |
|---|---|---|---|---|---|
| 1 | 1 | 1 | 2 | 3 | 1 |
| 3 | 3 | 1 | 2 | 2 | 3 |
| 4 | 4 | 4 | 4 | 4 | 4 |
| 1 | 2 | 1 | 2 | 2 | [1,2] |
Pandas 实现
直接通过apply逐行统计频次即可,代码如下:
import pandas as pd from collections import Counter # 读取/构造源数据 df = pd.DataFrame({ 'A': [1, 3, 4, 1], 'B': [1, 3, 4, 2], 'C': [1, 1, 4, 1], 'D': [2, 2, 4, 2] }) def row_stat(row): # 统计当前行四个列的值频次 count_res = Counter(row[['A', 'B', 'C', 'D']]) max_freq = max(count_res.values()) # 筛选所有达到最高频次的值 top_values = [val for val, freq in count_res.items() if freq == max_freq] # 单值返回标量,多值返回列表 return pd.Series([max_freq, top_values[0] if len(top_values) == 1 else top_values]) df[['#_identical', 'value']] = df.apply(row_stat, axis=1)
运行后输出结果和上述样例完全一致。
PySpark 实现
由于Spark为强类型引擎,同列无法同时存储标量和数组类型,生产环境建议统一将value列存为数组类型,单最高频值存长度为1的数组即可,若需完全对齐样例展示格式,可在最终输出层做格式转换。
实现代码如下:
from pyspark.sql import SparkSession from pyspark.sql.functions import udf from pyspark.sql.types import StructType, StructField, IntegerType, ArrayType from collections import Counter # 初始化Spark会话 spark = SparkSession.builder.appName("row_freq_stat").getOrCreate() # 构造/读取源数据 sdf = spark.createDataFrame([ (1, 1, 1, 2), (3, 3, 1, 2), (4, 4, 4, 4), (1, 2, 1, 2) ], schema=['A', 'B', 'C', 'D']) # 定义逐行统计UDF stat_schema = StructType([ StructField("#_identical", IntegerType(), False), StructField("value", ArrayType(IntegerType()), False) ]) @udf(stat_schema) def row_stat_spark(a, b, c, d): count_res = Counter([a, b, c, d]) max_freq = max(count_res.values()) top_values = [val for val, freq in count_res.items() if freq == max_freq] return max_freq, top_values # 计算新列 sdf = sdf.withColumn("stat_res", row_stat_spark("A", "B", "C", "D")) \ .select("A", "B", "C", "D", "stat_res.#_identical", "stat_res.value") sdf.show(truncate=False)
如果数据量较大,建议优先使用Spark内置高阶函数实现频次统计,避免Python UDF带来的序列化开销。
内容的提问来源于stack exchange,提问作者Karry
相关产品推荐
相关产品推荐

