如何在Spark Scala中仅当所有列B含mark时创建DataFrame
在Spark Scala中实现条件生成DataFrame的逻辑
假设我们有如下DataFrame:
| Column A | Column B |
|---|---|
| 123 | Mark 123 |
| 456 | Mark 456 |
| 789 | 789 |
要实现的逻辑:如果Column B的所有值都包含(不区分大小写的)"mark"字符,则生成仅保留Column A的DataFrame;否则返回空DataFrame。
实现代码
import org.apache.spark.sql.functions.{col, lower, expr} // 假设原始DataFrame名为df,SparkSession实例为spark // 方式1:统计不符合条件的记录数 val invalidRowCount = df.filter(!lower(col("Column B")).contains("mark")).count() val resultDF = if (invalidRowCount == 0) { df.select("Column A") } else { spark.emptyDataFrame } // 方式2:用every聚合函数直接判断所有行是否符合条件 val allRowsValid = df .agg(expr("every(lower(`Column B`) contains 'mark')").as("all_valid")) .first() .getBoolean(0) val resultDF = if (allRowsValid) { df.select("Column A") } else { spark.emptyDataFrame }
代码说明
- 使用
lower(col("Column B"))是为了忽略大小写匹配"mark",如果需要严格区分大小写,直接去掉lower函数即可。 - 方式1通过过滤出不包含"mark"的行并统计数量,若数量为0则说明所有行都符合要求;方式2利用Spark的
every聚合函数,直接判断所有行是否满足条件,代码更简洁。 - 列名包含空格,使用
expr时需要用反引号包裹列名,或者用col函数引用。
内容的提问来源于stack exchange,提问作者8919_racso
相关产品推荐
相关产品推荐

