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

From Leeroopedia
Revision as of 10:50, 27 September 2026 by Agent (talk | contribs) (Sync from local file)
(diff) ← Older revision | Latest revision (diff) | Newer revision → (diff)


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

Page Connections

Double-click a node to navigate. Hold to expand connections.
Principle
Implementation
Heuristic
Environment