PySpark计算逻辑实现:统计各Story_id对应的URI在主URL集合中的匹配数量
我来帮你解决这个PySpark的统计需求,思路其实很清晰——我们可以通过DataFrame的连接操作筛选出匹配的记录,再按story_id分组计数即可。下面是具体的方案和代码示例:
解决方案
方法1:内连接+分组计数(仅统计有匹配的story_id)
这种方法会只保留df1中uri存在于df2的记录,然后按story_id统计匹配数量,适合只关注有匹配结果的场景。
from pyspark.sql import SparkSession from pyspark.sql.functions import count, col # 初始化SparkSession spark = SparkSession.builder.appName("StoryUriMatchCount").getOrCreate() # 模拟示例输入数据(实际使用时替换为你的真实DataFrame) data1 = [ ("s1", "url1"), ("s1", "url2"), ("s1", "url3"), ("s2", "url2"), ("s2", "url4") ] df1 = spark.createDataFrame(data1, ["story_id", "uri"]) data2 = [("url1",), ("url3",), ("url4",)] df2 = spark.createDataFrame(data2, ["uri"]) # 核心逻辑:内连接筛选匹配记录,再分组计数 result_df = df1.join(df2, on="uri", how="inner") \ .groupBy("story_id") \ .agg(count("uri").alias("matching_uri_count")) # 查看结果 result_df.show()
执行后输出结果:
+---------+-------------------+ |story_id |matching_uri_count | +---------+-------------------+ |s1 |2 | |s2 |1 | +---------+-------------------+
方法2:左连接+分组计数(保留所有story_id)
如果需要统计所有story_id的匹配情况(包括没有匹配到df2中uri的,计数为0),可以使用左连接:
# 左连接保留所有story_id,统计非空uri的数量(即匹配次数) result_df_all = df1.join(df2, on="uri", how="left") \ .groupBy("story_id") \ .agg(count(col("uri")).alias("matching_uri_count")) result_df_all.show()
如果某个story_id没有匹配的uri,会显示计数为0,比如新增s3, url5的话,结果会包含s3, 0。
性能优化:使用广播变量
如果df2的主URL集合规模较小,建议使用broadcast广播df2到所有Executor节点,这样可以减少Shuffle操作,大幅提升大数据量下的处理性能:
from pyspark.sql.functions import broadcast # 广播df2后执行连接 result_df_optimized = df1.join(broadcast(df2), on="uri", how="inner") \ .groupBy("story_id") \ .agg(count("uri").alias("matching_uri_count")) result_df_optimized.show()
内容的提问来源于stack exchange,提问作者ilovupdates
相关产品推荐
相关产品推荐

