Spark DataFrame按换行拆分列并提取最后两项生成新列报错
解决Spark DataFrame拆分列并获取最后两项的问题
我明白你遇到的问题了——用Python列表的索引方式去操作Spark的Column对象肯定行不通,因为两者根本不是一回事儿!咱们一步步来搞定它。
为什么你的代码会返回None?
你写的df.withColumn('last_2', split(df.s, '\r\n')[-2])之所以失效,核心原因是:split()返回的是Spark Column类型的数组,不是普通的Python列表。Python的索引语法[-2]只能用来操作本地列表,对Spark Column完全不生效,自然就返回None了。
正确的解决方案
咱们分两种常用场景来处理,看你需要数组形式的最后两项,还是拼接成字符串的形式:
场景1:生成包含最后两项的数组列
用Spark的slice()函数,结合size()获取数组长度,就能精准截取最后两项:
from pyspark.sql import functions as F # 先拆分字符串为数组,再截取最后两项 df_result = df1.withColumn("split_array", F.split(F.col("s"), "\r\n")) \ .withColumn("last_2_items", F.slice(F.col("split_array"), F.size(F.col("split_array")) - 1, 2)) \ .drop("split_array") # 查看结果 df_result.show(truncate=False)
F.split(F.col("s"), "\r\n"):把s列按换行符拆分成数组F.size(F.col("split_array")):获取数组的总长度(你的示例里每行数组长度都是3)F.slice(..., size-1, 2):从数组的第size-1位(Spark数组索引是1-based)开始,截取2个元素,也就是最后两项
场景2:生成拼接后的字符串列
如果需要把最后两项合并成一个字符串(比如保留原换行符分隔),可以用element_at()(支持负索引,Spark 2.4+可用)+ concat_ws():
from pyspark.sql import functions as F df_result = df1.withColumn("split_array", F.split(F.col("s"), "\r\n")) \ .withColumn("last_2_str", F.concat_ws("\r\n", F.element_at(F.col("split_array"), -2), # 倒数第二项 F.element_at(F.col("split_array"), -1) # 倒数第一项 )) \ .drop("split_array") df_result.show(truncate=False)
兼容Spark 2.4以下版本(不支持负索引)
如果你的Spark版本较低,不能用负索引,就用size()计算正索引位置:
from pyspark.sql import functions as F df_result = df1.withColumn("split_array", F.split(F.col("s"), "\r\n")) \ .withColumn("last_2_str", F.concat_ws("\r\n", F.element_at(F.col("split_array"), F.size(F.col("split_array")) - 2), F.element_at(F.col("split_array"), F.size(F.col("split_array")) - 1) )) \ .drop("split_array")
完整可运行示例
把所有流程串起来,方便你直接测试:
from pyspark.sql import SparkSession from pyspark.sql import functions as F # 初始化SparkSession spark = SparkSession.builder.appName("SplitLastTwo").getOrCreate() # 创建示例DataFrame df1 = spark.createDataFrame([ ["001\r\nLuc Krier\r\n2363 Ryan Road, Long Lake South Dakota"], ["002\r\nJeanny Thorn\r\n2263 Patton Lane Raleigh North Carolina"], ["003\r\nTeddy E Beecher\r\n2839 Hartland Avenue Fond Du Lac Wisconsin"], ["004\r\nPhilippe Schauss\r\n1 Im Oberdorf Allemagne"], ["005\r\nMeindert I Tholen\r\nHagedoornweg 138 Amsterdam"] ]).toDF("s") # 执行处理(这里用拼接字符串的示例) df_result = df1.withColumn("split_array", F.split(F.col("s"), "\r\n")) \ .withColumn("last_2_str", F.concat_ws("\r\n", F.element_at(F.col("split_array"), -2), F.element_at(F.col("split_array"), -1) )) \ .drop("split_array") # 打印结果 df_result.show(truncate=False)
内容的提问来源于stack exchange,提问作者John Doe
相关产品推荐
相关产品推荐

