Implementation:Risingwavelabs Risingwave BulkListener
| Property | Value |
|---|---|
| File | java/connector-node/risingwave-sink-es-7/src/main/java/com/risingwave/connector/BulkListener.java
|
| Language | Java |
| Lines | 97 |
| Category | Listener |
| Package | com.risingwave.connector
|
Overview
BulkListener is a dual-purpose listener implementation that handles callbacks from both Elasticsearch and OpenSearch bulk request operations. It implements both org.elasticsearch.action.bulk.BulkProcessor.Listener and org.opensearch.action.bulk.BulkProcessor.Listener interfaces, providing a unified error tracking mechanism through the EsSink.RequestTracker. The listener logs bulk operation progress and records success or failure results for downstream error reporting.
Code Reference
Source Location
java/connector-node/risingwave-sink-es-7/src/main/java/com/risingwave/connector/BulkListener.java
Signature
public class BulkListener
implements org.elasticsearch.action.bulk.BulkProcessor.Listener,
org.opensearch.action.bulk.BulkProcessor.Listener {
public BulkListener(RequestTracker requestTracker);
}
Imports
import com.risingwave.connector.EsSink.RequestTracker;
import org.elasticsearch.action.bulk.BulkRequest;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
I/O Contract
Inputs
- Constructor: Takes a
RequestTrackerinstance for recording bulk operation outcomes.
Callback Methods (Elasticsearch)
- beforeBulk(long executionId, BulkRequest request): Logs the number of actions being sent to Elasticsearch.
- afterBulk(long executionId, BulkRequest request, BulkResponse response): On success, calls
requestTracker.addOkResult()with the action count. On failure, callsrequestTracker.addErrResult()with a formatted error message. - afterBulk(long executionId, BulkRequest request, Throwable failure): Called when the bulk operation raises a
Throwable; records the error viarequestTracker.addErrResult().
Callback Methods (OpenSearch)
The OpenSearch callback methods mirror the Elasticsearch ones exactly, operating on org.opensearch.action.bulk.BulkRequest and org.opensearch.action.bulk.BulkResponse types instead.
Outputs
The listener does not return values directly. It records results to the RequestTracker which is polled by the sink to detect errors and track completion.
Usage Examples
// Creating a BulkListener with a request tracker
RequestTracker requestTracker = new EsSink.RequestTracker();
BulkListener listener = new BulkListener(requestTracker);
// Used internally by BulkProcessor.Builder
BulkProcessor.builder(
(bulkRequest, bulkResponseActionListener) ->
client.bulkAsync(bulkRequest, RequestOptions.DEFAULT, bulkResponseActionListener),
listener
);
Related Pages
- ElasticBulkProcessorAdapter - Uses this listener for Elasticsearch bulk operations
- OpensearchBulkProcessorAdapter - Uses this listener for OpenSearch bulk operations
- BulkProcessorAdapter - Interface that the adapters implement