【问题标题】:How to filter a Spark RDD based on particular field value in Java?如何根据 Java 中的特定字段值过滤 Spark RDD?
【发布时间】:2015-10-07 12:52:02
【问题描述】:

我正在用 Java 创建一个 Spark 作业。这是我的代码。

我正在尝试从 CSV 文件中过滤记录。标头包含字段OIDCOUNTRY_NAME、......

不只是基于s.contains("CANADA") 进行过滤,我想更具体一些,比如我想根据COUNTRY_NAME.equals("CANADA") 进行过滤。 关于如何做到这一点的任何想法?

public static void main(String[] args) {
    String gaimFile = "hdfs://xx.yy.zz.com/sandbox/data/acc/mydata"; 

    SparkConf conf = new SparkConf().setAppName("Filter App");
    JavaSparkContext sc = new JavaSparkContext(conf);
    try{
        JavaRDD<String> gaimData = sc.textFile(gaimFile);

        JavaRDD<String> canadaOnly = gaimData.filter(new Function<String, Boolean>() {

            private static final long serialVersionUID = -4438640257249553509L;

            public Boolean call(String s) { 
               // My file id csv with header OID, COUNTRY_NAME, .....
               // here instead of just saying s.contains 
               // i would like to be more specific and say 
               // if COUNTRY_NAME.eqauls("CANADA)
               return s.contains("CANADA"); 
            }
        }); 

    }
    catch(Exception e){
        System.out.println("ERROR: G9 MatchUp Failed");
    }
    finally{
        sc.close();
    }
}

【问题讨论】:

    标签: java filter apache-spark


    【解决方案1】:

    您必须首先将您的值映射到自定义类:

    rdd.map(lines=>ConvertToCountry(line))
       .filter(country=>country == "CANADA")
    
    class Country{
      ...ctor that takes an array and fills properties...
      ...properties for each field from the csv...
    }
    
    ConvertToCountry(line: String){
      return new Country(line.split(','))
    }
    

    以上是 Scala 和伪代码的组合,但你应该明白这一点。

    【讨论】:

      猜你喜欢
      • 2020-08-22
      • 2014-11-20
      • 2021-04-10
      • 1970-01-01
      • 2017-02-20
      • 2016-02-29
      • 2021-11-03
      • 1970-01-01
      • 2017-10-22
      相关资源
      最近更新 更多