如何在PySpark中按行列索引为DataFrame列元素分配X.X格式序号
实现方法
你可以通过PySpark窗口函数生成从1开始的连续行号,再按规则拼接序号替换原有两个phrase字段即可,完整代码如下:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession(如果已有可以跳过该步骤) spark = SparkSession.builder.appName("assign_seq").getOrCreate() # 构造原始DataFrame df = spark.createDataFrame([ ('red apple', 'ripe banana', 0.5), ('late autumn', 'heavy rain', 0.1), ('speak loudly','quiet place', 0.9), ('extremely dangerous','fast running', 0.89) ], ["phrase1", "phrase2", 'common_persent']) # 定义窗口生成连续行号 window_spec = Window.orderBy(F.lit(1)) df = df.withColumn("row_id", F.row_number().over(window_spec)) # 按规则替换phrase字段值 df = df.withColumn("phrase1", F.concat(F.col("row_id"), F.lit(".1")).cast("double")) \ .withColumn("phrase2", F.concat(F.col("row_id"), F.lit(".2")).cast("double")) \ .drop("row_id") # 输出结果 df.show()
运行后输出的结果和你期望的完全一致:
+-------+-------+--------------+ |phrase1|phrase2|common_persent| +-------+-------+--------------+ | 1.1| 1.2| 0.5| | 2.1| 2.2| 0.1| | 3.1| 3.2| 0.9| | 4.1| 4.2| 0.89| +-------+-------+--------------+
注意事项
如果是大数据量场景,Spark分布式架构本身不保证数据读取的原始顺序,如果你需要严格稳定匹配输入时的行顺序,建议提前给数据添加可排序的业务字段(如自增ID、数据生成时间等),将上述窗口定义中的orderBy(F.lit(1))替换为对应业务字段即可。
内容的提问来源于stack exchange,提问作者Rory
相关产品推荐
相关产品推荐

