【问题标题】:Skip records in dataframe's map transformation在数据框的地图转换中跳过记录
【发布时间】:2021-06-04 16:19:41
【问题描述】:

我有一个 Spark 数据框,我在其上执行以下某些操作。我想知道如何跳过处理所有操作的某些记录

finalDf = df.map(mapFunc).reduceGroups(reduceFunc).map(mapFunc2).write().format().option().mode().save();

在 mapFunc 中,我想编写一个逻辑,如果某个条件为真,则不返回任何内容,并且考虑到该记录以进行进一步操作而必须退出 .reduceGroups(reduceFunc).map(mapFunc2).write().format().option().mode().save()

我尝试从地图返回 Optional.empty(),但代码在 reduceGroups 中失败并出现以下错误。

Exception in thread "main" java.util.NoSuchElementException: head of empty list
    at scala.collection.immutable.Nil$.head(List.scala:420)
    at scala.collection.immutable.Nil$.head(List.scala:417)
    at org.apache.spark.sql.catalyst.encoders.ExpressionEncoder$$anonfun$5.apply(ExpressionEncoder.scala:121)
    at org.apache.spark.sql.catalyst.encoders.ExpressionEncoder$$anonfun$5.apply(ExpressionEncoder.scala:120)
    at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
    at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
    at scala.collection.immutable.List.foreach(List.scala:381)
    at scala.collection.TraversableLike$class.map(TraversableLike.scala:234)
    at scala.collection.immutable.List.map(List.scala:285)
    at org.apache.spark.sql.catalyst.encoders.ExpressionEncoder$.tuple(ExpressionEncoder.scala:120)
    at org.apache.spark.sql.catalyst.encoders.ExpressionEncoder$.tuple(ExpressionEncoder.scala:187)
    at org.apache.spark.sql.expressions.ReduceAggregator.bufferEncoder(ReduceAggregator.scala:38)
    at org.apache.spark.sql.expressions.Aggregator.toColumn(Aggregator.scala:100)
    at org.apache.spark.sql.KeyValueGroupedDataset.reduceGroups(KeyValueGroupedDataset.scala:436)
    at org.apache.spark.sql.KeyValueGroupedDataset.reduceGroups(KeyValueGroupedDataset.scala:448)

数据框架构:

根 |-- id: 字符串(可为空=真) |-- 中:整数(可为空=真) |-- 响应:字符串(可为空=真) |-- 版本:字符串(可为空=真)

在 mapFunc 中,我通过执行子字符串并从中获取值来处理响应字符串。在某些情况下,响应字符串可以为空,因此在这种情况下,我不希望整个记录都在 finalDf 中;

Input

id mid responses version
A   1  "hello123" 1
B   1  "hello456" 2
A   2  "hello789" 5
A   1  "hello143" 4
B   3  "hello153" 6
C   3  ""         1

Output (Grouping by id, mid column)

id mid responses version
A   1  "143" 4
B   1  "456" 2
A   2  "789" 5
B   3  "153" 6

如果 (C,3) 的 (id,mid) 组合的响应不为空,那么它将在输出中。所以我想删除 mapFunc 中的 C,3。

【问题讨论】:

  • 什么都不返回是什么意思,你能不能也添加一些示例数据和预期的输出?
  • @koiralo:更新了问题
  • 你也可以添加mapFuc吗?或者您可以按原样返回空字符串并在下一步中使用过滤器

标签: apache-spark apache-spark-sql


【解决方案1】:

这里如果你有一个map 函数如下,那么你可以只返回同一行和filter 稍后使用过滤器

val myFunc = (row: Row) => {
  val number = row.getString(2).replaceAll("\\D+","")
  (row.getString(0), row.getInt(1), number, row.getInt(3))
}

df.map(myFunc)
  .toDF(df.columns: _*)
  .filter(trim($"responses") =!= "" )
  .show(false)

输出:

+---+---+---------+-------+
|id |mid|responses|version|
+---+---+---------+-------+
|A  |1  |123      |1      |
|B  |1  |456      |2      |
|A  |2  |789      |5      |
|A  |1  |143      |4      |
|B  |3  |153      |6      |
+---+---+---------+-------+

【讨论】:

    猜你喜欢
    • 2020-05-23
    • 1970-01-01
    • 1970-01-01
    • 2014-02-27
    • 1970-01-01
    • 2014-06-30
    • 1970-01-01
    • 2013-05-20
    • 1970-01-01
    相关资源
    最近更新 更多