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 BulkListener

From Leeroopedia


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 RequestTracker instance 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, calls requestTracker.addErrResult() with a formatted error message.
  • afterBulk(long executionId, BulkRequest request, Throwable failure): Called when the bulk operation raises a Throwable; records the error via requestTracker.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

Page Connections

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