【问题标题】:Cassandra With Spark connector - How Insert list Of Items to CassandraCassandra 与 Spark 连接器 - 如何将项目列表插入 Cassandra
【发布时间】:2021-12-10 02:56:51
【问题描述】:

使用 Java 使用 Cassandra 和 Spark 2.12 (3.2.0)。 Cassandra 连接器 3.1.0

我从 s3 获取的目的是进行预处理并将并行插入 Cassandra。

我遇到的问题是我对每个 s3 文件进行了预处理,其中包括要插入 Cassandra 的项目列表,如下所示:JavaRDD<List<SearchEntity>>

我应该如何将它传递给 cassandra(如代码示例中所示)?我看到它支持单个对象mapToRow。 也许我错过了什么?

使用以下代码

public static void main(String[] args) throws Exception {
    SparkConf conf = new SparkConf()
        .setAppName("Example Spark App")
        .setMaster("local[1]")
        .set("spark.cassandra.connection.host", "127.0.0.1");
        
    JavaSparkContext sparkContext = new JavaSparkContext(conf);
    sparkContext.hadoopConfiguration().set("fs.s3a.access.key", "XXXX");
    sparkContext.hadoopConfiguration().set("fs.s3a.secret.key", "YYYY");
    sparkContext.hadoopConfiguration().set("fs.s3a.endpoint", "XXXXX");
    sparkContext.hadoopConfiguration().set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem");
    sparkContext.hadoopConfiguration().set("mapreduce.input.fileinputformat.input.dir.recursive", "true");
                
    JavaPairRDD<String, PortableDataStream> javaPairRDD = sparkContext.binaryFiles("s3a://root/folder/");
        
    File ROOT = createTempFolder().getAbsoluteFile();
        
    JavaRDD<List<SearchEntity>> listJavaRDD = javaPairRDD.map(rdd -> {
            System.out.println("Working on TAR: " + rdd._1);
        
            DataInputStream stream = rdd._2.open();
        
            // some preprocess
            List<SearchEntity> toCassandraList = new WorkerTest(ROOT, stream).run();
        
            return toCassandraList;
        });
        
    // here I want to take List<SearchEntity> toCassandraList and save them
    // but I don't see how as it support only single object ..
    CassandraJavaUtil.javaFunctions(listJavaRDD)
        .writerBuilder("demoV2", "simple_search", 
                       CassandraJavaUtil.mapToRow(List<SearchEntity> list objects ...)) // here is problem
        .saveToCassandra();
        
    System.out.println("Finish run s3ToCassandra:");
    sparkContext.stop();
}

架构之前是手动配置的,仅用于测试目的。

CREATE TABLE simple_search (
    engine text,
    term text,
    time bigint,
    rank bigint,
    url text,
    domain text,
    pagenum bigint,
    descr text,
    display_url text,
    title text,
    type text,
    PRIMARY KEY ((engine, term), time , url, domain, pagenum)
) WITH CLUSTERING ORDER BY 
  (time DESC, url DESC,  domain DESC , pagenum DESC);

欢迎使用 Java 和 Scala 解决方案

【问题讨论】:

  • 什么是表架构?我不认为你有一个只包含条目列表的表格。另外,不使用 DataFrame API 而不是 RDD 的原因是什么?
  • 感谢@Alex Ott,这是二进制 tar 文件,我想用“/”从文件夹级别的 s3 加载它,没有找到如何使用 DataFrame 的任何示例,tar 也很大我需要在提取 tar 后合并一些内容文件并将它们转换为 cassandra 的 Pojo 列表(tar 架构在解压时是文件夹/文件层次结构)。关于模式和 KEYSPACE,我之前手动创建了它,但会在这里分享(顺便说一句,批量插入与简单的 cassandra java api 配合得很好),我想知道如何用 spark 来做到这一点。
  • 如果您可以将其转换为 DataFrame 以获得最佳实践,那就太好了:)
  • SearchEntity 按结构与表匹配?因为您需要做的第一件事是flatMap - 将RDD&lt;List&lt;SearchEntry&gt;&gt; 转换为RDD&lt;SearchEntry&gt;。另外,什么版本的 Spark Cassandra 连接器?
  • 是 SearchEntity 完全适合结构(它适用于基本的 java cassandra 连接器)。 3.1.0spark-cassandra-connector_2.12。可以分享一下怎么弄平吗?此外,它不会影响打开每个 JavaRDD - SearchEntry 的性能,因为现在我有数百万个 SearchEntry 项目......

标签: java apache-spark cassandra spark-cassandra-connector


【解决方案1】:

要写入数据,您需要在SearchEntity 上工作,而不是在SearchEntity 的列表上工作。为此,您需要使用flatMap 而不是普通的map

JavaRDD<SearchEntity> entriesRDD = javaPairRDD.flatMap(rdd -> {
        System.out.println("Working on TAR: " + rdd._1);
        DataInputStream stream = rdd._2.open();
        // some preprocess
        List<SearchEntity> toCassandraList = new WorkerTest(ROOT, stream).run();
        return toCassandraList;
    });

然后你就可以按照documentation

javaFunctions(rdd).writerBuilder("demoV2", "simple_search",
   mapToRow(SearchEntity.class)).saveToCassandra();

附:但请注意,如果您的 tar 太大,在创建 List&lt;SearchEntity&gt; 时可能会导致工作人员出现内存错误。根据 tar 文件中的文件格式,最好先解压缩数据,然后使用 Spark 读取它们。

【讨论】:

  • 解压磁盘上的数据?我有数百万的 tars 很多 IO,也许你熟悉一些 java/scala 库,我可以直接使用它从内存中的 tar.gz 按名称获取文件?
  • commons-compress 可能会有所帮助。我的评论很笼统——这完全取决于许多因素
  • 你的回答这不起作用,我得到了编译错误。我能做的是: JavaRDD rddList = javaPairRDD.map(rdd -> { DataInputStream stream = rdd._2.open(); System.out.println("Working on TAR:" + rdd._1); // 但我不确定 List listItem = new WorkerTest(ROOT, stream).run(); return listItem; }).flatMap(List::iterator);
  • 嗯,这可能是 Scala 和 Java 的区别。我的答案可以修改为返回 iterator()...
猜你喜欢
  • 2020-02-12
  • 2017-03-04
  • 2016-09-02
  • 1970-01-01
  • 2017-01-13
  • 2015-05-24
  • 2018-08-13
  • 2015-10-28
  • 1970-01-01
相关资源
最近更新 更多