【问题标题】:Scala - Using DF filter on multiple fieldsScala - 在多个字段上使用 DF 过滤器
【发布时间】:2018-12-01 05:53:34
【问题描述】:

我已经通过 stackoverflow 搜索了几天,但我只是没有找到以下问题的答案。我对 scala 编码真的很陌生,所以这可能是一个非常基本的问题。任何帮助将不胜感激。

我遇到的问题(出现错误)与最后一段代码有关。
我正在尝试从数据框中获取过滤后的记录子集,其中所有过滤后的记录都缺少一个或多个指定字段中的数据。

我在 Eclipse 中使用 Scala IDE Build 4.7.0。
我使用的 pom.xml 文件有 spark-core_2.11,版本 2.0.0

谢谢。
杰西

val source_path = args(0)
val source_file = args(1)

val vFile = sc.textFile(source_path + "/" + source_file)

val vSchema = StructType(
            StructField("FIELD_1",LongType,false)::
            StructField("FIELD_2",LongType,false)::
            StructField("FIELD_3",StringType,true)::
            StructField("FIELD_4",StringType,false)::
            StructField("FIELD_ADD_1",StringType,false)::
            StructField("FIELD_ADD_2",StringType,false)::
            StructField("FIELD_ADD_3",StringType,false)::
            StructField("FIELD_ADD_4",StringType,false)::
            StructField("FIELD_5",StringType,false)::
            StructField("FIELD_6",StringType,false)::
            StructField("FIELD_7",StringType,false)::
            StructField("FIELD_8",StringType,false)::
            Nil)

// val vRow = vFile.map(x=>x.split((char)30, -1)).map(x=> Row(
val vRow = vFile.map(x=>x.split("", -1)).map(x=> Row(
                            x(1).toLong,
                            x(2).toLong,
                            x(3).toString.trim(),
                            x(4).toString.trim(),
                            x(5).toString.trim(),
                            x(6).toString.trim(),
                            x(7).toString.trim(),
                            x(8).toString.trim(),
                            x(9).toString.trim(),
                            x(10).toString.trim(),
                            x(11).toString.trim(),
                            x(12).toString.trim()
                        ))

val dfData = sqlContext.createDataFrame(vRow.distinct(),vSchema)

val dfBlankRecords = dfData.filter(x => (
                    x.trim(col("FIELD_ADD_1")) == "" ||
                    x.trim(col("FIELD_ADD_2")) == "" ||
                    x.trim(col("FIELD_ADD_3")) == "" ||
                    x.trim(col("FIELD_ADD_4")) == ""
                ))

【问题讨论】:

  • 我会添加 apache-spark 标签以获得更好的可见性。在val vRow ... 行中,您正在执行x.split("", -1)。那里的空字符串是故意的吗?分割成单个字符的数组。
  • 另外,什么版本的火花?如果>= 1.6 有更好的方法将reading text files 直接导入数据集。
  • @TravisHegner,这些引号之间实际上有一个未打印的字符,作为文件中的列分隔符存在。我打算尝试使用以下行,以便更清晰一些,但我还没有弄清楚如何正确编写它。 val vSrcRow = vSrcFile.map(x=>x.split((char)30, -1)).map(x=> Row( 另外,根据我与 .scala 一起使用的 pom.xml 文件代码文件,我有 spark-core_2.11,版本 2.0.0。如果您愿意分享,我会很高兴看到将文本文件读入数据集的更好方法。

标签: eclipse scala dataframe filter


【解决方案1】:

spark.read.* 函数将数据直接读入 Dataset/Dataframe API,避免(在某种程度上)需要架构定义,并完全使用 RDD API。

val source_path = args(0)
val source_file = args(1)

val dfData = spark.read.textFile(source_path + "/" + source_file)
  .flatMap(l => {
    val a = l.split('\u001e'.toString, -1).map(_.trim())
    val f1 = a(0).toLong
    val f2 = a(1).toLong
    val Array(f3, f4, fa1, fa2, fa3, fa4, f5, f6, f7, f8) = a.slice(2,12)

    if (fa1 == "" ||
        fa2 == "" ||
        fa3 == "" ||
        fa4 == "") {
      Some(f1, f2, f3, f4, fa1, fa2, fa3, fa4, f5, f6, f7, f8)
    } else {
      None
    }
  }).toDF("FIELD_1", "FIELD_2", "FIELD_3", "FIELD_4",
          "FIELD_ADD_1", "FIELD_ADD_2", "FIELD_ADD_3", "FIELD_ADD_4",
          "FIELD_5", "FIELD_6", "FIELD_7", "FIELD_8")

我认为这会产生你想要的结果。我相信比我更好的人可以更好地优化它,并且使用更简洁的代码。

请注意,数组索引为零,如果您有意选择特定字段,则必须对其进行调整。我也不确定'\u001e'(十六进制值 30)是否是您拆分字符串所需的合适值。

【讨论】:

  • 看起来很棒。我将尝试将其合并到我的其余代码中,并让您知道它是如何工作的。我是否只需要定义长字段?对不起,我忘了指定角色的基础。我正在寻找 HEX 1E 或 ASCII 30。所以,看起来它也应该可以工作。非常感谢您的帮助!
  • 没问题。 .split() 已返回 Array[String],因此无需转换字符串字段。
  • 非常感谢您提供的信息。我已经适应了我正在尝试做的事情。但是,我尝试使用的文件相当大(它有 268 列)。因此,当我将所有这些都列在 Some() 方法中时,我收到以下错误:“too many arguments for method apply: (x: A)Some[A] in object Some”。因为,我不熟悉 Some() 方法,所以我不确定如何解决这个问题。该方法的参数限制是多少?
  • Some() 方法实际上会生成一个Option,它基本上是一个长度为 1 的Iterable。在执行.flatMap() 时,选项会变平。所以参数实际上表示为Some[TupleX],但是TupleX 的长度限制为22 项(阅读:Tuple2-Tuple22)。
  • 不幸的是,这个解决方案根本不适用于那么多列。是否可以将其余列保留为Strings?如果是这样,您可以发出Some[(Long, Long, Array[String])],这实际上是一个更清洁的解决方案。如果没有,您拥有的所有列类型是什么?
猜你喜欢
  • 2021-09-07
  • 1970-01-01
  • 2011-01-07
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多