【发布时间】: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<List<SearchEntry>>转换为RDD<SearchEntry>。另外,什么版本的 Spark Cassandra 连接器? -
是 SearchEntity 完全适合结构(它适用于基本的 java cassandra 连接器)。
3.1.0 spark-cassandra-connector_2.12 。可以分享一下怎么弄平吗?此外,它不会影响打开每个 JavaRDD - SearchEntry 的性能,因为现在我有数百万个 SearchEntry 项目......
标签: java apache-spark cassandra spark-cassandra-connector