PySpark中按空格拆分字符串并转为单词数组的实现方案
解决PySpark中拆分查询字符串为去重单词数组的问题
要实现将每个AnonID对应的搜索查询字符串拆分为独立单词并去重的需求,你可以通过两种高效的PySpark内置函数方案实现,无需使用UDF:
方案一:拆分后展开再聚合去重
这种方式逻辑直观,先将每个查询字符串拆分为单词数组,再展开为单个单词行,最后按AnonID收集去重后的单词:
import pyspark.sql.functions as F from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() # 读取原始数据 df = spark.read.option("header", "true") \ .option("delimiter", "\t") \ .option("inferSchema", "true") \ .csv("/content/drive/MyDrive/my_data.txt") # 拆分每个Query为单词数组,并展开为每行一个单词 df_words = df.withColumn("word", F.explode(F.split("Query", "\\s+"))) # 按AnonID分组,收集去重后的单词数组 result = df_words.groupBy("AnonID").agg(F.collect_set("word").alias("New_Data")) # 查看结果 result.show()
方案二:聚合阶段直接处理数组(更高效)
这种方式无需生成中间行,直接在聚合时合并所有单词数组并去重,减少数据shuffle开销:
import pyspark.sql.functions as F from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() # 读取原始数据 df = spark.read.option("header", "true") \ .option("delimiter", "\t") \ .option("inferSchema", "true") \ .csv("/content/drive/MyDrive/my_data.txt") # 分组后直接拆分、合并数组并去重 result = df.groupBy("AnonID").agg( F.array_distinct(F.flatten(F.collect_list(F.split("Query", "\\s+")))).alias("New_Data") ) # 查看结果 result.show()
关键说明:
- 使用
F.split("Query", "\\s+")而非简单空格拆分,可匹配多个连续空格,避免生成空字符串单词。 collect_set会自动去重并返回数组,而array_distinct+flatten组合则是先合并所有单词数组,再去重,两种方式都能得到你需要的结果。
内容的提问来源于stack exchange,提问作者Amit BenDavid
相关产品推荐
相关产品推荐

