You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 01:10:28