【问题标题】:how to filter few rows in a table using Scala如何使用Scala过滤表中的几行
【发布时间】:2019-08-22 05:14:05
【问题描述】:

使用 Scala: 我有一个emp表如下

id, name,   dept,   address
1,  a,  10, hyd
2,  b,  10, blr
3,  a,  5,  chn
4,  d,  2,  hyd
5,  a,  3,  blr
6,  b,  2,  hyd

代码:

val inputFile = sc.textFile("hdfs:/user/edu/emp.txt"); 
val inputRdd = inputFile.map(iLine => (iLine.split(",")(0),
                             iLine.split(",")(1), 
                             iLine.split(",")(3)
                            )); 
// filtering only few columns Now i want to pull hyd addressed employees complete data 

问题:我不想打印所有 emp 详细信息,我只想打印几个都来自 hyd 的 emp 详细信息。

  1. 我已将此 emp 数据集加载到 Rdd 中
  2. 我用 ',' 分割了这个 Rdd
  3. 现在我只想打印 hyd 寻址的 emp。

【问题讨论】:

  • 代码在哪里?
  • val inputFile = sc.textFile("hdfs:/user/edu/emp.txt"); val inputRdd = inputFile.map(iLine => (iLine.split(",")(0),iLine.split(",")(1), iLine.split(",")(3))); // 只过滤几列现在我想提取 hyd 寻址员工的完整数据
  • 你在 RDD 上尝试过过滤方法吗?
  • 请分享代码以过滤 rdd 上的行

标签: scala apache-spark


【解决方案1】:

我认为以下解决方案将有助于解决您的问题。

  val fileName = "/path/stact_test.txt"
  val strRdd = sc.textFile(fileName).map { line =>
    val data = line.split(",")
    (data(0), data(1), data(3))
  }.filter(rec=>rec._3.toLowerCase.trim.equals("hyd"))

拆分数据后,使用元组 RDD 中的第 3 项过滤位置。

输出:

(1,  a, hyd)
(4,  d,  hyd)
(6,  b,  hyd)

【讨论】:

    【解决方案2】:

    您可以尝试使用数据框

    
    val viewsDF=spark.read.text("hdfs:/user/edu/emp.txt")
    val splitedViewsDF = viewsDF.withColumn("id", split($"value",",").getItem(0))
                                .withColumn("name", split($"value", ",").getItem(1))
                                .withColumn("address", split($"value", ",").getItem(3))
                                .drop($"value")
                                .filter(df("address").equals("hyd") )
    
    

    【讨论】:

    • thnx 但是,我收到错误::25: 错误:使用替代方法重载方法值过滤器:(func: org.apache.spark.api.java.function.FilterFunction[org .apache.spark.sql.Row])org.apache.spark.sql.Dataset[org.apache.spark.sql.Row] (func: org.apache.spark.sql.Row => Boolean)org .apache.spark.sql.Dataset[org.apache.spark.sql.Row] (conditionExpr: String)org.apache.spark.sql.Dataset[org.apache.spark.sql.Row] (条件:org.apache.spark.sql.Column)org.apache.spark.sql.Dataset[org.apache.spark.sql.Row]不能应用于(布尔)
    • 试试 df("address").equals("hyd")
    猜你喜欢
    • 1970-01-01
    • 2016-02-05
    • 2012-10-18
    • 2019-12-19
    • 1970-01-01
    • 1970-01-01
    • 2021-02-09
    • 1970-01-01
    • 2016-09-14
    相关资源
    最近更新 更多