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 ElasticBulkProcessorAdapter: Difference between revisions

From Leeroopedia
Auto-imported from implementations/Risingwavelabs_Risingwave_ElasticBulkProcessorAdapter.md
 
Sync from local file
 
Line 111: Line 111:
== Related Pages ==
== Related Pages ==


* [[Risingwavelabs_Risingwave_BulkProcessorAdapter|BulkProcessorAdapter]] - Interface implemented by this class
* [[Implementation:Risingwavelabs_Risingwave_BulkProcessorAdapter|BulkProcessorAdapter]] - Interface implemented by this class
* [[Risingwavelabs_Risingwave_ElasticRestHighLevelClientAdapter|ElasticRestHighLevelClientAdapter]] - Client adapter used for async bulk operations
* [[Implementation:Risingwavelabs_Risingwave_ElasticRestHighLevelClientAdapter|ElasticRestHighLevelClientAdapter]] - Client adapter used for async bulk operations
* [[Risingwavelabs_Risingwave_BulkListener|BulkListener]] - Listener handling bulk operation callbacks
* [[Implementation:Risingwavelabs_Risingwave_BulkListener|BulkListener]] - Listener handling bulk operation callbacks
* [[Risingwavelabs_Risingwave_EsSinkConfig|EsSinkConfig]] - Configuration driving bulk processor settings
* [[Implementation:Risingwavelabs_Risingwave_EsSinkConfig|EsSinkConfig]] - Configuration driving bulk processor settings
* [[Risingwavelabs_Risingwave_OpensearchBulkProcessorAdapter|OpensearchBulkProcessorAdapter]] - OpenSearch counterpart of this adapter
* [[Implementation:Risingwavelabs_Risingwave_OpensearchBulkProcessorAdapter|OpensearchBulkProcessorAdapter]] - OpenSearch counterpart of this adapter


[[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/ElasticBulkProcessorAdapter.java
Language Java
Lines 104
Category Adapter
Package com.risingwave.connector

Overview

ElasticBulkProcessorAdapter is the Elasticsearch 7-specific implementation of the BulkProcessorAdapter interface. It wraps the Elasticsearch BulkProcessor API to translate generic bulk add/delete operations into Elasticsearch-native UpdateRequest and DeleteRequest calls. The adapter configures the bulk processor with batch size limits, flush intervals, concurrent request settings, and exponential backoff retry policies based on values from EsSinkConfig.

Code Reference

Source Location

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

Signature

public class ElasticBulkProcessorAdapter implements BulkProcessorAdapter {
    public ElasticBulkProcessorAdapter(
            RequestTracker requestTracker,
            ElasticRestHighLevelClientAdapter client,
            EsSinkConfig config);
    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 com.risingwave.connector.EsSink.RequestTracker;
import java.util.concurrent.TimeUnit;
import org.elasticsearch.action.bulk.BackoffPolicy;
import org.elasticsearch.action.bulk.BulkProcessor;
import org.elasticsearch.action.delete.DeleteRequest;
import org.elasticsearch.action.update.UpdateRequest;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.common.unit.ByteSizeUnit;
import org.elasticsearch.common.unit.ByteSizeValue;
import org.elasticsearch.core.TimeValue;
import org.elasticsearch.xcontent.XContentType;

I/O Contract

Constructor

Initializes the Elasticsearch BulkProcessor with the following configuration from EsSinkConfig:

Setting Config Method Description
Bulk actions getBatchNumMessages() Number of actions before automatic flush
Bulk size getBatchSizeKb() Size in KB before automatic flush
Flush interval hardcoded 5 seconds
Concurrent requests getConcurrentRequests() Number of concurrent bulk requests
Backoff policy hardcoded Exponential backoff starting at 100ms, up to 3 retries
Retry on conflict getRetryOnConflict() Number of retries on version conflicts

addRow

Creates an Elasticsearch UpdateRequest with docAsUpsert(true), meaning documents are inserted if they do not exist and updated if they do. The document type is set to "_doc". Optional routing is applied if provided.

deleteRow

Creates an Elasticsearch DeleteRequest targeting the specified index and key with document type "_doc". Optional routing is applied if provided.

Both addRow and deleteRow call requestTracker.addWriteTask() before adding the request to the bulk processor.

Usage Examples

// Create the adapter with an Elasticsearch client
RequestTracker tracker = new EsSink.RequestTracker();
ElasticRestHighLevelClientAdapter client = new ElasticRestHighLevelClientAdapter(host, config);
ElasticBulkProcessorAdapter adapter = new ElasticBulkProcessorAdapter(tracker, client, config);

// Upsert a document
adapter.addRow("my_index", "key1", "{\"name\": \"value\"}", null);

// Delete a document with routing
adapter.deleteRow("my_index", "key2", "routing_value");

// Flush and close
adapter.flush();
adapter.awaitClose(30, TimeUnit.SECONDS);

Related Pages