You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用Pyspark/Pandas逐行统计多列相同值数量及对应取值

逐行多列频次统计实现

你可以根据使用的工具栈选择Pandas或PySpark实现对应逻辑:

规则说明

源数据固定包含A、B、C、D四个数值列,逐行计算后输出两个新列:

  • #_identical:当前行四个列值中,出现次数最高的频次
  • value:对应最高频次的取值,仅一个最高频值时存该值,多个值频次并列最高时存所有值的列表
    计算逻辑参考样例:
ABCD#_identicalvalue
111231
331223
444444
12122[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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 06:12:13