Spark 1.6 DataFrame提取p、s列值相同且w列共有的行
解决Spark 1.6 DataFrame筛选共有(p,s)组合的问题
我来帮你搞定这个需求!要实现保留所有(p,s)组合为所有w共有的行,可以按以下步骤用Spark 1.6的API来完成:
步骤1:明确核心逻辑
我们需要找出那些每个w都至少包含一次的(p,s)组合,然后保留原DataFrame中属于这些组合的所有行。
步骤2:具体实现代码
首先假设你已经有了目标DataFrame(我先帮你构造示例数据方便测试):
import org.apache.spark.sql.Row import org.apache.spark.sql.types._ // 定义Schema val schema = StructType(Seq( StructField("w", StringType, nullable = false), StructField("p", StringType, nullable = false), StructField("s", IntegerType, nullable = false) )) // 示例数据 val data = Seq( Row("w1", "p1", 0), Row("w1", "p1", 1), Row("w1", "p1", 2), Row("w1", "p2", 0), Row("w1", "p2", 1), Row("w2", "p1", 0), Row("w2", "p1", 1), Row("w2", "p3", 0), Row("w2", "p3", 1), Row("w2", "p3", 2), Row("w2", "p4", 0), Row("w1", "p4", 1), Row("w3", "p1", 0), Row("w3", "p1", 1), Row("w3", "p2", 0), Row("w3", "p3", 0), Row("w3", "p4", 0), Row("w4", "p1", 0), Row("w4", "p1", 1) ) // 创建DataFrame val df = sqlContext.createDataFrame(sc.parallelize(data), schema)
接下来执行核心转换:
// 1. 统计总共有多少个不同的w val totalUniqueWs = df.select("w").distinct().count() // 2. 找出所有被所有w都包含的(p,s)组合 val sharedPsPairs = df.select("p", "s", "w") .distinct() // 每个w对同一个(p,s)只计一次 .groupBy("p", "s") .count() .filter(s"count = $totalUniqueWs") // 只保留被所有w覆盖的(p,s) .select("p", "s") // 只保留(p,s)列用于后续过滤 // 3. 过滤原DataFrame,只保留符合条件的行 val resultDf = df.join(sharedPsPairs, Seq("p", "s"), "inner")
步骤3:查看结果
执行以下代码查看最终输出:
resultDf.orderBy("w", "p", "s").show()
输出结果会和你期望的一致,只保留(p1,0)和(p1,1)这两个所有w都共有的组合对应的行:
+---+---+---+ | p| s| w| +---+---+---+ | p1| 0| w1| | p1| 1| w1| | p1| 0| w2| | p1| 1| w2| | p1| 0| w3| | p1| 1| w3| | p1| 0| w4| | p1| 1| w4| +---+---+---+
逻辑说明
- 第一步统计总w数:确保我们知道一个(p,s)需要被多少个不同的w包含才算“共有”
- 第二步去重+分组计数:避免同一个w多次统计同一个(p,s),然后筛选出覆盖所有w的组合
- 第三步内连接:只保留原DataFrame中属于这些共有(p,s)组合的行
内容的提问来源于stack exchange,提问作者mbaxi
相关产品推荐
相关产品推荐

