Jump to content

Connect SuperML | Leeroopedia MCP: Equip your AI agents with best practices, code verification, and debugging knowledge. Powered by Leeroo — building Organizational Superintelligence. Contact us at founders@leeroo.com.

Implementation:Risingwavelabs Risingwave BulkProcessorAdapter: Difference between revisions

From Leeroopedia
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

java/connector-node/risingwave-sink-es-7/src/main/java/com/risingwave/connector/BulkProcessorAdapter.java

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

  • addRow and deleteRow may throw InterruptedException if the calling thread is interrupted while waiting.
  • awaitClose may throw InterruptedException if 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