You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Java Client Api(DMSDK)从MarkLogic生成130万条数据CSV的问询

Using MarkLogic DMSDK to Generate 1.3M-Record CSV (No Load Balancer Required)

Got it, let's tackle this—DMSDK is perfect here because it natively distributes work across your cluster nodes without needing a load balancer, which should cut down on that 1+ hour runtime you saw with CORB. Here's a step-by-step guide with code examples to get you up and running efficiently:

1. Set Up Dependencies

First, add the DMSDK dependency to your project. If using Maven, include this in your pom.xml:

<dependency>
    <groupId>com.marklogic</groupId>
    <artifactId>marklogic-data-movement</artifactId>
    <version>6.3.0</version> <!-- Use the latest version compatible with your MarkLogic server -->
</dependency>

2. Initialize the Data Movement Manager

Connect directly to all your cluster nodes (no load balancer needed—DMSDK handles node distribution automatically). Create a DataMovementManager with credentials for every node:

import com.marklogic.client.DatabaseClient;
import com.marklogic.client.DatabaseClientFactory;
import com.marklogic.client.datamovement.DataMovementManager;
import com.marklogic.client.datamovement.DataMovementManagerFactory;
import java.util.Arrays;
import java.util.List;

// Create clients for each cluster node
List<DatabaseClient> clients = Arrays.asList(
    DatabaseClientFactory.newClient("node1.your-domain.com", 8000, "your-username", "your-password", DatabaseClientFactory.Authentication.DIGEST),
    DatabaseClientFactory.newClient("node2.your-domain.com", 8000, "your-username", "your-password", DatabaseClientFactory.Authentication.DIGEST),
    // Add all remaining cluster nodes here
);

// Initialize DMSDK manager with all node clients
DataMovementManager dmm = DataMovementManagerFactory.newInstance(clients);

3. Configure the Query Batcher

Define which documents you want to export. Use a QueryBatcher to target your 1.3M records—you can use a structured query, XQuery, or URI pattern. For example, if your records live in a specific collection:

import com.marklogic.client.datamovement.QueryBatcher;
import com.marklogic.client.query.StructuredQueryBuilder;
import com.marklogic.client.query.StructuredQueryDefinition;

StructuredQueryBuilder sqb = new StructuredQueryBuilder();
StructuredQueryDefinition query = sqb.collection("target-record-collection");

// Create a batcher that processes documents in batches
QueryBatcher batcher = dmm.newQueryBatcher(query)
    .withBatchSize(2000) // Adjust based on server resources—start with 1k-5k
    .withThreadCount(10); // Match to your cluster's capacity (e.g., 2-3 threads per node)

4. Add a Document Processor to Extract CSV Fields

Create a DocumentProcessor to extract needed fields from each document and format them as CSV rows. To optimize performance, use server-side XQuery to extract fields before sending data to your client (reduces network transfer):

import com.marklogic.client.datamovement.DocumentProcessor;
import com.marklogic.client.document.JSONDocumentManager;
import com.marklogic.client.io.StringHandle;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;

// Thread-safe list to collect CSV rows (or write directly to file with synchronization)
List<String> csvRows = Collections.synchronizedList(new ArrayList<>());
csvRows.add("field1,field2,field3"); // Add CSV header

DocumentProcessor csvExtractor = new DocumentProcessor() {
    @Override
    public void processDocument(DatabaseClient client, String uri) {
        JSONDocumentManager docMgr = client.newJSONDocumentManager();
        // Server-side XQuery to extract only required fields (faster than fetching full documents)
        String xquery = "declare variable $uri as xs:string external; " +
            "let $doc := fn:doc($uri) " +
            "return concat( " +
                "fn:encode-for-uri($doc//field1/text()), ',', " +
                "fn:encode-for-uri($doc//field2/text()), ',', " +
                "fn:encode-for-uri($doc//field3/text()) " +
            ")";
        
        StringHandle resultHandle = new StringHandle();
        client.newServerEval()
            .xquery(xquery)
            .addVariable("uri", uri)
            .eval(resultHandle);
        
        csvRows.add(resultHandle.get());
    }
};

batcher.onUrisReady(csvExtractor);

5. Handle Batch Completion and Write to File

Add a BatchListener to write batches of rows to your CSV file once processed—this avoids holding all 1.3M rows in memory:

import com.marklogic.client.datamovement.BatchListener;
import com.marklogic.client.datamovement.QueryBatch;
import java.io.BufferedWriter;
import java.io.FileWriter;
import java.io.IOException;

BufferedWriter writer = new BufferedWriter(new FileWriter("output.csv"));
writer.write("field1,field2,field3\n"); // Write header first

batcher.onBatchSuccess(new BatchListener<QueryBatch>() {
    @Override
    public void processEvent(QueryBatch batch) {
        // Write all rows from this batch to the file
        synchronized (writer) {
            try {
                for (int i = batch.getStartIndex(); i < batch.getEndIndex(); i++) {
                    writer.write(csvRows.get(i) + "\n");
                }
                writer.flush();
            } catch (IOException e) {
                throw new RuntimeException("Failed to write batch to CSV", e);
            }
        }
    }
});

// Handle batch failures gracefully
batcher.onBatchFailure((batch, throwable) -> {
    System.err.println("Failed to process batch starting at index " + batch.getStartIndex());
    throwable.printStackTrace();
});

6. Run the Job and Clean Up

Start the batcher, wait for completion, then close all resources:

import com.marklogic.client.datamovement.JobTicket;

JobTicket ticket = dmm.startJob(batcher);
dmm.waitForCompletion(ticket);

// Clean up resources
writer.close();
for (DatabaseClient client : clients) {
    client.release();
}
dmm.release();

Key Optimization Tips

  • Batch Size & Thread Count: Adjust withBatchSize and withThreadCount based on your cluster's CPU/memory. Too large a batch can cause memory issues; too small wastes overhead.
  • Server-Side Extraction: Always use server-side XQuery/JavaScript to extract fields instead of fetching full documents—this cuts down network traffic drastically.
  • Avoid In-Memory Storage: Use batch writing to the file instead of holding all rows in memory (critical for 1.3M records).
  • Monitor Progress: Add a ProgressListener to track processing status:
    batcher.onProgressUpdate(progress -> {
        System.out.println("Processed " + progress.getSuccessEvents() + " records so far");
    });
    

This approach should be significantly faster than CORB because DMSDK distributes work directly across all cluster nodes without a load balancer, and you have fine-grained control over processing logic to optimize performance.

内容的提问来源于stack exchange,提问作者Private

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.19 04:06:08