如何在PostgreSQL和/或PySpark中移除text_1列中与text_2列单词完全重复或作为其子串的单词
用PostgreSQL和PySpark实现文本过滤需求
没问题,我来帮你搞定这个文本处理的需求——移除text_1列中那些在text_2列里完全重复,或是作为子串存在的单词,最后合并剩余text_1和所有text_2的内容。下面分别给出两种工具的实现方式:
PostgreSQL 实现方式
PostgreSQL可以通过子查询判断匹配条件,再用字符串聚合来生成最终结果。假设你的数据存在名为your_table的表中:
步骤1:过滤符合移除条件的text_1单词
先筛选出需要保留的text_1单词:
SELECT text_1 FROM your_table t1 WHERE NOT EXISTS ( SELECT 1 FROM your_table t2 -- 判断当前text_1单词是否是text_2的完全匹配项或子串 WHERE t2.text_2 = t1.text_1 OR t2.text_2 LIKE '%' || t1.text_1 || '%' );
步骤2:合并剩余text_1与所有text_2内容
把过滤后的text_1和所有text_2的单词拼接成一行输出:
SELECT -- 拼接保留的text_1单词 string_agg(filtered_text_1, ' ') || ' ' || -- 拼接所有text_2单词 string_agg(text_2, ' ') AS final_result FROM ( -- 子查询获取过滤后的text_1 SELECT text_1 AS filtered_text_1 FROM your_table t1 WHERE NOT EXISTS ( SELECT 1 FROM your_table t2 WHERE t2.text_2 = t1.text_1 OR t2.text_2 LIKE '%' || t1.text_1 || '%' ) ) filtered_t1, ( -- 子查询获取所有text_2 SELECT text_2 FROM your_table ) all_t2;
执行这段SQL后,就能得到你想要的结果(注:原需求输出里的text_1 text_2应为表头说明,实际数据部分就是过滤后拼接的单词串)。
PySpark 实现方式
PySpark可以通过收集text_2的单词集合,再用UDF(用户自定义函数)过滤text_1的内容,最后拼接成字符串。
步骤1:初始化SparkSession并加载数据
from pyspark.sql import SparkSession from pyspark.sql.functions import col, collect_set, array_join, udf, array_filter from pyspark.sql.types import ArrayType, StringType, BooleanType # 初始化SparkSession spark = SparkSession.builder.appName("TextFilterDemo").getOrCreate() # 加载你的数据 data = [("astro", "lumen"), ("cosm", "planet"), ("microcosm", "astronomy"), ("planet", "magnitude")] df = spark.createDataFrame(data, ["text_1", "text_2"])
步骤2:收集text_2的单词集合并广播
为了提高效率,我们把text_2的所有单词收集到一个集合,并用广播变量分发到各个节点:
# 收集text_2的所有唯一单词 text_2_words = df.select(collect_set("text_2")).first()[0] # 广播变量,避免每个任务重复传输数据 broadcast_text2 = spark.sparkContext.broadcast(text_2_words)
步骤3:定义过滤函数并处理数据
用UDF过滤text_1的单词,再拼接成最终结果:
# 定义过滤函数:判断单词是否需要保留 def should_keep_word(word): t2_words = broadcast_text2.value # 如果单词在text_2中完全匹配,或是text_2某个单词的子串,就返回False(不保留) return not any(word == t2 or word in t2 for t2 in t2_words) # 注册UDF keep_word_udf = udf(should_keep_word, BooleanType()) # 收集所有text_1和text_2为数组,过滤后拼接 result_df = df.agg( collect_set("text_1").alias("text1_array"), collect_set("text_2").alias("text2_array") ).withColumn( "filtered_text1", # 过滤text1数组中需要移除的单词 array_filter(col("text1_array"), lambda x: keep_word_udf(x)) ).withColumn( "final_result", # 拼接过滤后的text1和所有text2 array_join(col("filtered_text1"), " ") + " " + array_join(col("text2_array"), " ") ) # 查看结果 result_df.select("final_result").show(truncate=False)
运行这段代码后,输出的final_result列就是你需要的结果。
内容的提问来源于stack exchange,提问作者red_quark
相关产品推荐
相关产品推荐

