【问题标题】:How to write JavaRDD to marklogic database如何将 JavaRDD 写入 marklogic 数据库
【发布时间】: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 任务不可序列化。

【问题讨论】:

标签: apache-spark marklogic marklogic-8


【解决方案1】:

将数据输入 MarkLogic 的最简单方法是通过 HTTP 和客户端 REST API - 特别是 /v1/documents 端点 - http://docs.marklogic.com/REST/client/management

有多种方法可以优化这一点,例如通过写入集,但根据您的问题,我认为首先要决定的是 - 您想为每条记录编写什么样的文档?您的示例在 CSV 中显示了 5 列 - 通常,您将编写包含 5 个字段/元素的 JSON 或 XML 文档,每个字段/元素均基于列索引命名。因此,您需要编写一些代码来生成该 JSON/XML,然后使用您喜欢的任何 HTTP 客户端(一种选择是 MarkLogic Java 客户端 API)将该文档写入 MarkLogic。

这解决了您如何将 JavaRDD 写入 MarkLogic 的问题 - 但如果您的目标是尽快将数据从 CSV 获取到 MarkLogic,那么请跳过 Spark 并使用 mlcp - https://docs.marklogic.com/guide/mlcp/import#id_70366 - 这涉及零编码。

【讨论】:

  • 我已经更新了我的问题,添加了 java 客户端 api 的代码以写入 marklogic 数据库。我正在将记录写为 XML 数据。但得到异常:sparkexception 任务不可序列化。无法找出异常的原因。如果您有任何线索,请告诉我。此外,如果我写入数据库的方式不正确,那么建议我使用 THE ONE。
  • Java 客户端 API 中有多种方法可用于编写文档或一组文档,对于本示例,“writeAs”就可以了。
  • 为了将问题隔离到 DatabaseClient 不可序列化,我建议不要使用 JAXBContext,而只需做一些非常简单的事情,比如创建一个 XML 字符串,然后通过 StringHandle 编写它。如果可行,那么您就知道问题出在 JAXB 上,并且某些类不是可序列化的。如果这不起作用,那么问题出在 DatabaseClient 上,我认为 Sam Mefford 的以下建议与对 DatabaseClient 的静态引用是一个很好的尝试。
【解决方案2】:

修改了spark streaming guide的例子,在这里你必须实现连接和编写特定于数据库的逻辑。

public void send(JavaRDD<String> rdd) {
    rdd.foreachPartition(new VoidFunction<Iterator<String>>() {
      @Override
      public void call(Iterator<String> partitionOfRecords) {
        // ConnectionPool is a static, lazily initialized pool of
        Connection connection = ConnectionPool.getConnection();
        while (partitionOfRecords.hasNext()) {
          connection.send(partitionOfRecords.next());
        }
        ConnectionPool.returnConnection(connection); // return to the pool
        // for future reuse
      }
    });
  }

【讨论】:

  • 您的回答并没有明确说明我如何将数据写入 marklogic 数据库。你能详细说明你的答案吗?
  • 如何为marklogic数据库创建connectionPool?
  • 请参阅docs.marklogic.com/guide/java/intro 它有关于写入 MarkLogic 的示例。希望对您有所帮助。
  • 我知道marklogic java api可以写入marklogic。我想要做的是使用 spark api 将数据写入 marklogic。您提供的示例,如果这有助于将数据写入 marklogic 数据库,那么我的第一个问题将是如何创建示例中特定于 marklogic 的连接对象。
  • Connection 和 ConnectionPool 只是示例类(而不是具体类或 Spark API 的一部分)。您必须编写自己的具有类似功能的类。
【解决方案3】:

我想知道您是否只需要确保您在 VoidFunction 内部访问且在其外部实例化的所有内容都是可序列化的(请参阅this page)。 DatabaseClient 和 XMLDocumentManager 当然不可序列化,因为它们是连接的资源。但是,不要在 VoidFunction 中实例化 DatabaseClient 是对的,因为那样效率会降低(尽管它会起作用)。我不知道以下想法是否适用于火花。但我猜你可以创建一个持有单例 DatabaseClient 实例的类:

public static class MLClient {
  private static DatabaseClient singleton;
  private MLClient() {}

  public static DatabaseClient get(DatabaseClientFactory.Bean connectionInfo) {
    if ( connectionInfo == null ) {
      throw new IllegalArgumentException("connectionInfo cannot be null");
    }
    if ( singleton == null ) {
      singleton = connectionInfo.newClient();
    }
    return singleton;
  }
}

然后您只需在您的 VoidFunction 之外创建一个可序列化的 DatabaseClientFactory.Bean,这样您的身份验证信息仍然是集中的

DatabaseClientFactory.Bean connectionInfo = 
  new DatabaseClientFactory.Bean();
connectionInfo.setHost("localhost");
connectionInfo.setPort(8000);
connectionInfo.setUser("admin");
connectionInfo.setPassword("admin");
connectionInfo.setAuthenticationValue("digest");

然后在您的 VoidFunction 中,您可以像这样获得单例 DatabaseClient 和新的 XMLDocumentManager:

DatabaseClient client = MLClient.get(connectionInfo);
XMLDocumentManager docMgr = client.newXMLDocumentManager();

【讨论】:

  • 感谢 Sam 的回答。我知道​​序列化数据库客户端是不可取的,但它是 spark 的 VoidFunction,它需要所有资源都必须序列化。实际上我正在寻找带有 spark 的 marklogic 连接器,但没有找到任何连接器,因此尝试使用 spark 的 java 客户端 api 本身。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2011-07-02
  • 2020-04-20
  • 1970-01-01
  • 1970-01-01
  • 2022-11-23
  • 1970-01-01
  • 2023-02-25
相关资源
最近更新 更多