【发布时间】:2017-05-13 13:27:59
【问题描述】:
我正在使用 marklogic 数据库评估 spark。我已经阅读了一个 csv 文件,现在我有一个 JavaRDD 对象,我必须将其转储到 marklogic 数据库中。
SparkConf conf = new SparkConf().setAppName("org.sparkexample.Dataload").setMaster("local");
JavaSparkContext sc = new JavaSparkContext(conf);
JavaRDD<String> data = sc.textFile("/root/ml/workArea/data.csv");
SQLContext sqlContext = new SQLContext(sc);
JavaRDD<Record> rdd_records = data.map(
new Function<String, Record>() {
public Record call(String line) throws Exception {
String[] fields = line.split(",");
Record sd = new Record(fields[0], fields[1], fields[2], fields[3],fields[4]);
return sd;
}
});
我想将这个 JavaRDD 对象写入 marklogic 数据库。
是否有任何 spark api 可用于更快地写入 marklogic 数据库?
比方说,如果我们不能将 JavaRDD 直接写入 marklogic,那么实现这一点的正确方法是什么?
这是我用来将 JavaRDD 数据写入 marklogic 数据库的代码,如果这样做是错误的方法,请告诉我。
final DatabaseClient client = DatabaseClientFactory.newClient("localhost",8070, "MLTest");
final XMLDocumentManager docMgr = client.newXMLDocumentManager();
rdd_records.foreachPartition(new VoidFunction<Iterator<Record>>() {
public void call(Iterator<Record> partitionOfRecords) {
while (partitionOfRecords.hasNext()) {
Record record = partitionOfRecords.next();
System.out.println("partitionOfRecords - "+record.toString());
String docId = "/example/"+record.getID()+".xml";
JAXBContext context = JAXBContext.newInstance(Record.class);
JAXBHandle<Record> handle = new JAXBHandle<Record>(context);
handle.set(record);
docMgr.writeAs(docId, handle);
}
}
});
client.release();
我已经使用 java 客户端 api 来编写数据,但是即使 POJO 类 Record 正在实现 Serializable 接口,我也遇到了异常。请让我知道可能是什么原因以及如何解决。
org.apache.spark.sparkexception 任务不可序列化。
【问题讨论】:
-
我对 Spark 了解不多,但这里有一些资源可能会有所帮助:How to use MarkLogic in Apache Spark applications; Putting Spark to Work With MarkLogic.
-
谢谢戴夫,我已经看到了那个例子。该示例将输出数据存储到 hdfs 而不是 marklogic 数据库。
标签: apache-spark marklogic marklogic-8