基于Java Client Api(DMSDK)从MarkLogic生成130万条数据CSV的问询
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
withBatchSizeandwithThreadCountbased 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
ProgressListenerto 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

