PySpark DataFrame是否存在类似pandas的iloc、cumsum的对应功能方法?
PySpark 支持实现类似 pandas 中iloc和cumsum的功能,具体实现方法如下:
1 逐行累积求和(等价pandas的
cumsum(axis=1),适配链梯法三角表计算场景) PySpark 没有直接提供行方向累积求和的内置方法,可以通过数组聚合逻辑实现,示例代码如下:
from pyspark.sql import functions as F # 按你三角表的列顺序整理需要计算累积和的列名列表,可根据实际列名调整过滤规则 cum_cols = [c for c in ft.columns if c != "事故年标识列"] # 实现行方向累积求和 ft = ft.select( "*", F.aggregate( F.array(*cum_cols), F.expr("array<double>()"), lambda acc, current: F.array_union( acc, F.array(F.coalesce(current, F.lit(0)) + F.coalesce(F.element_at(acc, F.size(acc)), F.lit(0))) ) ).alias("cumulative_array") ) # 将累积和数组拆回原有列 for idx, col_name in enumerate(cum_cols): ft = ft.withColumn(col_name, F.element_at("cumulative_array", idx + 1)) ft = ft.drop("cumulative_array")
如果需要列方向的累积求和,等价于 pandas 的cumsum(axis=0),直接用窗口函数即可:
from pyspark.sql.window import Window # 注意指定明确的排序列,保证累积顺序符合预期 window_spec = Window.orderBy("你指定的排序列").rowsBetween(Window.unboundedPreceding, 0) ft = ft.withColumn("列方向累积和列名", F.sum("目标计算列").over(window_spec))
2 位置索引(等价pandas的
iloc功能) PySpark 是分布式计算框架,没有默认固定的行顺序,所以实现iloc前需要先给数据加全局有序的行索引:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 加行索引,索引从0开始和pandas对齐,注意排序列要能保证全局顺序唯一 window_spec = Window.orderBy("你指定的全局排序列") ft = ft.withColumn("row_index", F.row_number().over(window_spec) - 1)
按行位置取数据(等价df.iloc[行索引])
# 取索引为2的行数据 target_row = ft.filter(F.col("row_index") == 2).collect()
按行列位置取指定值(等价df.iloc[行索引, 列索引])
col_list = ft.columns # 取第2行、第3列的值(行索引1,列索引2) target_value = ft.filter(F.col("row_index") == 1).select(col_list[2]).collect()[0][0]
注意:如果你需要和 pandas 计算结果完全对齐,一定要指定唯一且稳定的排序规则,避免分布式调度导致行顺序变化,引发计算结果偏差。
内容的提问来源于stack exchange,提问作者srivats29
相关产品推荐
相关产品推荐

