代码之家  ›  专栏  ›  技术社区  ›  Stepan Yakovenko

spark mlib ALS算法的Maven依赖地狱[重复]

  •  0
  • Stepan Yakovenko  · 技术社区  · 7 年前

    我有一小段java代码来获得apache spark推荐:

    公共静态类评级实现可序列化{ 私人电影ID;

        public Rating() {}
    
        public Rating(int userId, int movieId, float rating, long timestamp) {
            this.userId = userId;
            this.movieId = movieId;
            this.rating = rating;
            this.timestamp = timestamp;
        }
    
        public int getUserId() {
            return userId;
        }
    
        public int getMovieId() {
            return movieId;
        }
    
        public float getRating() {
            return rating;
        }
    
        public long getTimestamp() {
            return timestamp;
        }
    
        public static Rating parseRating(String str) {
            String[] fields = str.split(",");
            if (fields.length != 4) {
                throw new IllegalArgumentException("Each line must contain 4 fields");
            }
            int userId = Integer.parseInt(fields[0]);
            int movieId = Integer.parseInt(fields[1]);
            float rating = Float.parseFloat(fields[2]);
            long timestamp = Long.parseLong(fields[3]);
            return new Rating(userId, movieId, rating, timestamp);
        }
    }
    
    static String parse(String str) {
        Pattern pat = Pattern.compile("\\[[0-9.]*,[0-9.]*]");
        Matcher matcher = pat.matcher(str);
        int count = 0;
        StringBuilder sb = new StringBuilder();
        while (matcher.find()) {
            count++;
            String substring = str.substring(matcher.start(), matcher.end());
            String itstr = substring.split(",")[0].substring(1);
            sb.append(itstr + " ");
        }
        return sb.toString().trim();
    }
    
    static TreeMap<Long, String> res = new TreeMap<>();
    
    public static void add(long k, String v) {
        res.put(k, v);
    }
    
    public static void main(String[] args) throws IOException {
        Logger.getLogger("org").setLevel(Level.OFF);
        Logger.getLogger("akka").setLevel(Level.OFF);
        SparkSession spark = SparkSession
                .builder()
                .appName("SomeAppName")
                .config("spark.master", "local")
                .getOrCreate();
        JavaRDD<Rating> ratingsRDD = spark
                .read().textFile(args[0]).javaRDD()
                .map(Rating::parseRating);
        Dataset<Row> ratings = spark.createDataFrame(ratingsRDD, Rating.class);
        ALS als = new ALS()
                .setMaxIter(1)
                .setRegParam(0.01)
                .setUserCol("userId")
                .setItemCol("movieId")
                .setRatingCol("rating");
        ALSModel model = als.fit(ratings);
        model.setColdStartStrategy("drop");
        Dataset<Row> rowDataset = model.recommendForAllUsers(50);
        rowDataset.foreach((ForeachFunction<Row>) row -> {
            String str = row.toString();
            long l = Long.parseLong(str.substring(1).split(",")[0]);
            add(l, parse(str));
        });
        BufferedWriter bw = new BufferedWriter(new FileWriter(args[1]));
        for (long l = 0; l < res.lastKey(); l++) {
            if (!res.containsKey(l)) {
                bw.write("\n");
                continue;
            }
            String str = res.get(l);
            bw.write(str);
        }
        bw.close();
    }
    

    }

    我尝试在pom.xml中使用不同的依赖项来运行它,但所有变体都失败了。这个:

        <dependency>
            <groupId>com.sparkjava</groupId>
            <artifactId>spark-core</artifactId>
            <version>2.8.0</version>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-mllib_2.12</artifactId>
            <version>2.4.0</version>
        </dependency>
        <dependency>
            <groupId>org.slf4j</groupId>
            <artifactId>slf4j-api</artifactId>
            <version>1.6.4</version>
        </dependency>
    
        <dependency>
            <groupId>org.slf4j</groupId>
            <artifactId>slf4j-log4j12</artifactId>
            <version>1.6.4</version>
        </dependency>
    

    ,为了修复它,我添加了

    spark-sql-kafka-0-10_2.10 2.0.2

    ClassNotFoundException:org.apache.spark.internal.Logging$class ,要修复它,我添加了另一个:

        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming_2.11</artifactId>
            <version>2.2.2</version>
        </dependency>
    
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-core_2.10</artifactId>
            <version>2.2.2</version>
        </dependency>
    
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-sql_2.10</artifactId>
            <version>2.2.2</version>
        </dependency>
    
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming-kafka-0-8_2.11</artifactId>
            <version>2.2.2</version>
        </dependency>
    

    现在它失败了 java.lang.NoClassDefFoundError:scala/collection/GenTraversableOnce 为了解决这个问题,我尝试了十几种其他组合,都失败了,最后一种是

        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming_2.11</artifactId>
            <version>2.2.2</version>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-mllib_2.11</artifactId>
            <version>2.2.2</version>
        </dependency>
    
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-core_2.11</artifactId>
            <version>2.2.2</version>
        </dependency>
    
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-sql_2.11</artifactId>
            <version>2.2.2</version>
        </dependency>
    
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming-kafka-0-8_2.11</artifactId>
            <version>2.2.2</version>
        </dependency>
    

    这又让我

    <dependencies>
        <dependency> <!-- Spark dependency -->
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-core_2.11</artifactId>
            <version>2.0.1</version>
        </dependency>
        <dependency> <!-- Spark dependency -->
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-mllib_2.11</artifactId>
            <version>2.0.1</version>
        </dependency>
        <dependency> <!-- Spark dependency -->
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-sql_2.11</artifactId>
            <version>2.0.1</version>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming_2.11</artifactId>
            <version>2.0.1</version>
        </dependency>
        <dependency>
            <groupId>org.apache.bahir</groupId>
            <artifactId>spark-streaming-twitter_2.11</artifactId>
            <version>2.0.1</version>
        </dependency>
    </dependencies>
    

    (这仍然让我感到惊讶 java.lang.ClassNotFoundException:text.DefaultSource ))

    我还尝试了在这个问题中发布的依赖项,但它们也失败了: Resolving dependency problems in Apache Spark

    https://github.com/stiv-yakovenko/sparkrec

    1 回复  |  直到 7 年前
        1
  •  0
  •   Stepan Yakovenko    7 年前

    <dependencies>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-core_2.11</artifactId>
            <version>2.4.0</version>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-sql_2.11</artifactId>
            <version>2.4.0</version>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-mllib_2.11</artifactId>
            <version>2.4.0</version>
        </dependency>
        <dependency>
            <groupId>org.scala-lang</groupId>
            <artifactId>scala-library</artifactId>
            <version>2.11.8</version>
        </dependency>
    </dependencies>
    

    您必须使用这些确切的版本,否则它将以多种方式崩溃。

    推荐文章