PySpark中基于索引偏移列值的实现(禁用UDF)
PySpark 列偏移实现方案
需求分析
根据INDEX字段的值,将每行的A/B列数据向左偏移对应数量的位置(每一组A_i/B_i为一个偏移单位),空出的右侧列填充空值,且禁止使用UDF。
实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when # 初始化SparkSession spark = SparkSession.builder.appName("ColumnShift").getOrCreate() # 构建初始数据 data = [ (0, "00a", "00b", "01a", "01b", "02a", "02b", "03a", "03b"), (1, None, None, "11a", "11b", "12a", "12b", "13a", "13b"), (2, None, None, None, None, "21a", "22b", "23a", "23b"), (3, None, None, None, None, None, None, "33a", "33b") ] columns = ["INDEX", "A_0", "B_0", "A_1", "B_1", "A_2", "B_2", "A_3", "B_3"] df = spark.createDataFrame(data, schema=columns) # 定义列偏移逻辑 shifted_df = df.select( col("INDEX"), # 处理A_0列:根据INDEX取对应位置的A列 when(col("INDEX") == 0, col("A_0")) .when(col("INDEX") == 1, col("A_1")) .when(col("INDEX") == 2, col("A_2")) .when(col("INDEX") == 3, col("A_3")) .alias("A_0"), # 处理B_0列 when(col("INDEX") == 0, col("B_0")) .when(col("INDEX") == 1, col("B_1")) .when(col("INDEX") == 2, col("B_2")) .when(col("INDEX") == 3, col("B_3")) .alias("B_0"), # 处理A_1列 when(col("INDEX") == 0, col("A_1")) .when(col("INDEX") == 1, col("A_2")) .when(col("INDEX") == 2, col("A_3")) .alias("A_1"), # 处理B_1列 when(col("INDEX") == 0, col("B_1")) .when(col("INDEX") == 1, col("B_2")) .when(col("INDEX") == 2, col("B_3")) .alias("B_1"), # 处理A_2列 when(col("INDEX") == 0, col("A_2")) .when(col("INDEX") == 1, col("A_3")) .alias("A_2"), # 处理B_2列 when(col("INDEX") == 0, col("B_2")) .when(col("INDEX") == 1, col("B_3")) .alias("B_2"), # 处理A_3列 when(col("INDEX") == 0, col("A_3")).alias("A_3"), # 处理B_3列 when(col("INDEX") == 0, col("B_3")).alias("B_3") ) # 查看结果 shifted_df.show()
逻辑说明
- 针对每一个目标列(如
A_0、B_0),通过when函数根据INDEX的值匹配对应的源列:- 当
INDEX=0时,保持原列值不变; - 当
INDEX=1时,取右侧一组列(A_1/B_1)的值填充到当前列; - 以此类推,
INDEX=N时,取右侧第N组列的值填充到当前列; - 没有匹配条件的场景自动填充空值。
- 当
- 全程使用PySpark内置函数实现,未使用UDF,符合要求。
内容的提问来源于stack exchange,提问作者Cri
相关产品推荐
相关产品推荐

