如何在PySpark数据帧中按±2小时间隔过滤重复西班牙语行
PySpark 过滤符合特定条件的西班牙语行
我有一个PySpark数据帧,需要过滤掉符合特定条件的西班牙语行。具体逻辑为:若某条西班牙语行存在另一条与之日期(date)、标签(tags)、来源(source)完全相同,且时间间隔在±2小时内的记录,则移除该西班牙语行。
测试数据
test_data = [('1', '2022-09-01' , '07:30:29' , '[tech, fx]' , 'YouTube' , 'english' ,'some text here'), ('2', '2022-09-01' , '07:30:29' , '[finance, fx]' , 'YouTube' , 'english' ,'some text here'), ('3', '2022-09-02' , '06:30:29' , '[tech, banking]' , 'YouTube' , 'english' ,'some text here'), ('4', '2022-09-02' , '07:20:29' , '[tech, banking]' , 'YouTube' , 'spanish' ,'Spanish Text'), ('5', '2022-09-03' , '07:12:55' , '[finance, fx]' , 'YouTube' , 'english' ,'some text here'), ('6', '2022-09-05' , '09:12:55' , '[computer]' , 'Instagram' , 'spanish' ,'Spanish Text'),] test_data = spark.sparkContext.parallelize(test_data).toDF(['id', 'date', 'time', 'tags', 'source', 'language', 'text'])
数据帧初始状态:
+---+----------+--------+---------------+---------+--------+--------------+ | id| date| time| tags| source|language| text| +---+----------+--------+---------------+---------+--------+--------------+ | 1|2022-09-01|07:30:29| [tech, fx]| YouTube| english|some text here| | 2|2022-09-01|07:30:29| [finance, fx]| YouTube| english|some text here| | 3|2022-09-02|06:30:29|[tech, banking]| YouTube| english|some text here| | 4|2022-09-02|07:20:29|[tech, banking]| YouTube| spanish| Spanish Text| | 5|2022-09-03|07:12:55| [finance, fx]| YouTube| english|some text here| | 6|2022-09-05|09:12:55| [computer]|Instagram| spanish| Spanish Text| +---+----------+--------+---------------+---------+--------+--------------+
示例说明:在该示例中,仅需移除第4行。
解决方案
步骤1:合并日期与时间为时间戳
首先将date和time字段合并为标准时间戳,方便后续计算时间间隔:
from pyspark.sql import functions as F df_with_ts = test_data.withColumn( "timestamp", F.to_timestamp(F.concat(F.col("date"), F.lit(" "), F.col("time")), "yyyy-MM-dd HH:mm:ss") )
步骤2:自关联匹配符合条件的记录
通过自关联,找到每个西班牙语行对应的、满足date/tags/source一致且时间在±2小时内的其他记录:
joined_df = df_with_ts.alias("a").join( df_with_ts.alias("b"), (F.col("a.date") == F.col("b.date")) & (F.col("a.tags") == F.col("b.tags")) & (F.col("a.source") == F.col("b.source")) & (F.col("a.id") != F.col("b.id")) & # 排除自身匹配 (F.abs(F.unix_timestamp("a.timestamp") - F.unix_timestamp("b.timestamp")) <= 7200), # 2小时=7200秒 how="left" )
步骤3:过滤保留目标行
保留所有非西班牙语行,以及没有匹配到符合条件记录的西班牙语行:
filtered_df = joined_df.filter( (F.col("a.language") != "spanish") | (F.col("b.id").isNull()) ).select("a.*").drop("timestamp")
验证结果
执行filtered_df.show()后,输出如下:
+---+----------+--------+---------------+---------+--------+--------------+ | id| date| time| tags| source|language| text| +---+----------+--------+---------------+---------+--------+--------------+ | 1|2022-09-01|07:30:29| [tech, fx]| YouTube| english|some text here| | 2|2022-09-01|07:30:29| [finance, fx]| YouTube| english|some text here| | 3|2022-09-02|06:30:29|[tech, banking]| YouTube| english|some text here| | 5|2022-09-03|07:12:55| [finance, fx]| YouTube| english|some text here| | 6|2022-09-05|09:12:55| [computer]|Instagram| spanish| Spanish Text| +---+----------+--------+---------------+---------+--------+--------------+
内容的提问来源于stack exchange,提问作者user19495470
相关产品推荐
相关产品推荐

