Implementation:Risingwavelabs Risingwave BulkProcessorAdapter: Difference between revisions
Appearance
Auto-imported from implementations/Risingwavelabs_Risingwave_BulkProcessorAdapter.md |
Sync from local file |
||
| Line 86: | Line 86: | ||
== Related Pages == | == Related Pages == | ||
* [[Risingwavelabs_Risingwave_ElasticBulkProcessorAdapter|ElasticBulkProcessorAdapter]] - Elasticsearch implementation of this interface | * [[Implementation:Risingwavelabs_Risingwave_ElasticBulkProcessorAdapter|ElasticBulkProcessorAdapter]] - Elasticsearch implementation of this interface | ||
* [[Risingwavelabs_Risingwave_OpensearchBulkProcessorAdapter|OpensearchBulkProcessorAdapter]] - OpenSearch implementation of this interface | * [[Implementation:Risingwavelabs_Risingwave_OpensearchBulkProcessorAdapter|OpensearchBulkProcessorAdapter]] - OpenSearch implementation of this interface | ||
* [[Risingwavelabs_Risingwave_BulkListener|BulkListener]] - Listener used by both adapter implementations | * [[Implementation:Risingwavelabs_Risingwave_BulkListener|BulkListener]] - Listener used by both adapter implementations | ||
* [[Risingwavelabs_Risingwave_EsSinkConfig|EsSinkConfig]] - Configuration used to initialize the adapters | * [[Implementation:Risingwavelabs_Risingwave_EsSinkConfig|EsSinkConfig]] - Configuration used to initialize the adapters | ||
[[Category:Implementations]] | [[Category:Implementations]] | ||
[[Category:Implementations]] | [[Category:Implementations]] | ||
Latest revision as of 10:50, 27 September 2026
| 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