Implementation:Risingwavelabs Risingwave BulkProcessorAdapter
Appearance
| Property | Value |
|---|---|
| File | java/connector-node/risingwave-sink-es-7/src/main/java/com/risingwave/connector/BulkProcessorAdapter.java
|
| Language | Java |
| Lines | 30 |
| Category | Interface |
| Package | com.risingwave.connector
|
Overview
BulkProcessorAdapter is an interface that defines the contract for bulk processor adapters supporting both Elasticsearch and OpenSearch. It provides an abstraction layer over the different client APIs for bulk write operations, enabling the ES sink connector to work with either backend through a unified interface. The interface declares methods for adding rows (upserts), deleting rows, flushing pending operations, and awaiting graceful shutdown.
Code Reference
Source Location
Signature
public interface BulkProcessorAdapter {
public void addRow(String index, String key, String doc, String routing) throws InterruptedException;
public void deleteRow(String index, String key, String routing) throws InterruptedException;
public void flush();
public void awaitClose(long timeout, TimeUnit unit) throws InterruptedException;
}
Imports
import java.util.concurrent.TimeUnit;
I/O Contract
Methods
| Method | Parameters | Description |
|---|---|---|
| addRow | index, key, doc (JSON), routing |
Adds a document upsert request to the bulk processor. The routing parameter is optional (may be null).
|
| deleteRow | index, key, routing |
Adds a document delete request to the bulk processor. The routing parameter is optional (may be null).
|
| flush | (none) | Forces the bulk processor to send all pending requests immediately. |
| awaitClose | timeout, unit |
Waits for the bulk processor to complete all pending operations and shut down within the specified timeout. |
Exceptions
addRowanddeleteRowmay throwInterruptedExceptionif the calling thread is interrupted while waiting.awaitClosemay throwInterruptedExceptionif interrupted during the shutdown wait.
Usage Examples
// Using the adapter to write and delete documents
BulkProcessorAdapter adapter = ...; // ElasticBulkProcessorAdapter or OpensearchBulkProcessorAdapter
// Upsert a document
adapter.addRow("my_index", "doc_key_1", "{\"field\": \"value\"}", null);
// Delete a document
adapter.deleteRow("my_index", "doc_key_2", null);
// Flush pending operations
adapter.flush();
// Graceful shutdown
adapter.awaitClose(30, TimeUnit.SECONDS);
Related Pages
- ElasticBulkProcessorAdapter - Elasticsearch implementation of this interface
- OpensearchBulkProcessorAdapter - OpenSearch implementation of this interface
- BulkListener - Listener used by both adapter implementations
- EsSinkConfig - Configuration used to initialize the adapters
Page Connections
Double-click a node to navigate. Hold to expand connections.
Principle
Implementation
Heuristic
Environment