【问题标题】:Read data saved by spark redis using Java使用Java读取spark redis保存的数据
【发布时间】:2020-12-20 11:20:48
【问题描述】:

我使用spark-redis 将数据集保存到 Redis。 然后我使用Spring data redis读取这些数据:

这个对象我保存到redis:

@Getter
@Setter
@AllArgsConstructor
@NoArgsConstructor
@Builder
@RedisHash("collaborative_filtering")
public class RatingResult implements Serializable {
    private static final long serialVersionUID = 8755574422193819444L;

    @Id
    private String id;

    @Indexed
    private int user;

    @Indexed
    private String product;

    private double productN;
    private double rating;
    private float prediction;

    public static RatingResult convert(Row row) {
        int user = row.getAs("user");
        String product = row.getAs("product");
        double productN = row.getAs("productN");
        double rating = row.getAs("rating");
        float prediction = row.getAs("prediction");
        String id = user + product;

        return RatingResult.builder().id(id).user(user).product(product).productN(productN).rating(rating)
                .prediction(prediction).build();
    }

}

使用 spark-redis 保存对象:

JavaRDD<RatingResult> result = ...
...
sparkSession.createDataFrame(result, RatingResult.class).write().format("org.apache.spark.sql.redis")
            .option("table", "collaborative_filtering").mode(SaveMode.Overwrite).save();

存储库:

@Repository
public interface RatingResultRepository extends JpaRepository<RatingResult, String> {

}

我无法读取使用 Spring data redis 保存在 Redis 中的数据,因为 spark-redis 和 spring data redis 保存的结构数据不一样(我检查了 spark-redis 和 spring data redis 创建的键的值是通过使用命令不同:redis-cli -p 6379 keys \*redis-cli hgetall $key)

那么如何读取已经使用 Java 或任何 Java 库保存的数据?

【问题讨论】:

    标签: java apache-spark redis spring-data-redis spark-redis


    【解决方案1】:

    以下对我有用。

    从 spark-redis 写入数据。

    我在这里使用 Scala,但它与您在 Java 中使用的基本相同。我唯一改变的是我添加了一个.option("key.column", "id") 来指定哈希ID。

        val ratingResult = new RatingResult("1", 1, "product1", 2.0, 3.0, 4)
    
        val result: JavaRDD[RatingResult] = spark.sparkContext.parallelize(Seq(ratingResult)).toJavaRDD()
        spark
          .createDataFrame(result, classOf[RatingResult])
          .write
          .format("org.apache.spark.sql.redis")
          .option("key.column", "id")
          .option("table", "collaborative_filtering")
          .mode(SaveMode.Overwrite)
          .save()
    

    在 spring-data-redis 中我有以下内容:

    @Getter
    @Setter
    @AllArgsConstructor
    @NoArgsConstructor
    @Builder
    @RedisHash("collaborative_filtering")
    public class RatingResult implements Serializable {
        private static final long serialVersionUID = 8755574422193819444L;
    
        @Id
        private String id;
    
        @Indexed
        private int user;
    
        @Indexed
        private String product;
    
        private double productN;
        private double rating;
        private float prediction;
    
        @Override
        public String toString() {
            return "RatingResult{" +
                    "id='" + id + '\'' +
                    ", user=" + user +
                    ", product='" + product + '\'' +
                    ", productN=" + productN +
                    ", rating=" + rating +
                    ", prediction=" + prediction +
                    '}';
        }
    }
    

    我使用 CrudRepository 而不是 JPA:

    @Repository
    public interface RatingResultRepository extends CrudRepository<RatingResult, String> {
    
    }
    

    查询:

         RatingResult found = ratingResultRepository.findById("1").get();
         System.out.println("found = " + found);
    

    输出:

    found = RatingResult{id='null', user=1, product='product1', productN=2.0, rating=3.0, prediction=4.0}
    

    您可能会注意到 id 字段未填充,因为存储的 spark-redis 具有哈希 id 而不是哈希属性。

    【讨论】:

    • 感谢您的回答,我已经通过id搜索成功了,但是有些情况下,我想搜索其他字段(用户,...)。我该如何实施?
    • 我找到了解决方案:查找所有键-> 按键查找RatingResult -> 其他字段过滤:RedisConnection redisConnection = null; try { redisConnection = redisTemplate.getConnectionFactory().getConnection(); ScanOptions options = ScanOptions.scanOptions().match("collaborative_filtering*") .count(100).build(); Cursor c = redisConnection.scan(options); while (c.hasNext()) { String id = new String((byte[]) c.next()); //TODO Find by id and filter } } finally { redisConnection.close(); }
    猜你喜欢
    • 2019-09-24
    • 1970-01-01
    • 2020-08-24
    • 1970-01-01
    • 2019-03-14
    • 1970-01-01
    • 2018-03-30
    • 2018-01-24
    • 1970-01-01
    相关资源
    最近更新 更多