This is an automated email from the ASF dual-hosted git repository.

guoweijie pushed a commit to branch main
in repository 
https://gitbox.apache.org/repos/asf/flink-connector-elasticsearch.git


The following commit(s) were added to refs/heads/main by this push:
     new da2ef1f  [FLINK-35504] Improve Elasticsearch 8 connector observability
da2ef1f is described below

commit da2ef1fa6d5edd3cf1328b11632929fd2c99f567
Author: Mingliang Liu <[email protected]>
AuthorDate: Sat Jun 1 16:06:32 2024 -0700

    [FLINK-35504] Improve Elasticsearch 8 connector observability
---
 .../sink/Elasticsearch8AsyncWriter.java            | 23 +++++++++++++++++++---
 1 file changed, 20 insertions(+), 3 deletions(-)

diff --git 
a/flink-connector-elasticsearch8/src/main/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriter.java
 
b/flink-connector-elasticsearch8/src/main/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriter.java
index a3eb8df..1163bf2 100644
--- 
a/flink-connector-elasticsearch8/src/main/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriter.java
+++ 
b/flink-connector-elasticsearch8/src/main/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriter.java
@@ -62,6 +62,13 @@ public class Elasticsearch8AsyncWriter<InputT> extends 
AsyncSinkWriter<InputT, O
     private boolean close = false;
 
     private final Counter numRecordsOutErrorsCounter;
+    /**
+     * A counter to track number of records that are returned by Elasticsearch 
as failed and then
+     * retried by this writer.
+     */
+    private final Counter numRecordsSendPartialFailureCounter;
+    /** A counter to track the number of bulk requests that are sent to 
Elasticsearch. */
+    private final Counter numRequestSubmittedCounter;
 
     private static final FatalExceptionClassifier 
ELASTICSEARCH_FATAL_EXCEPTION_CLASSIFIER =
             FatalExceptionClassifier.createChain(
@@ -103,11 +110,15 @@ public class Elasticsearch8AsyncWriter<InputT> extends 
AsyncSinkWriter<InputT, O
         checkNotNull(metricGroup);
 
         this.numRecordsOutErrorsCounter = 
metricGroup.getNumRecordsOutErrorsCounter();
+        this.numRecordsSendPartialFailureCounter =
+                metricGroup.counter("numRecordsSendPartialFailure");
+        this.numRequestSubmittedCounter = 
metricGroup.counter("numRequestSubmitted");
     }
 
     @Override
     protected void submitRequestEntries(
             List<Operation> requestEntries, Consumer<List<Operation>> 
requestResult) {
+        numRequestSubmittedCounter.inc();
         LOG.debug("submitRequestEntries with {} items", requestEntries.size());
 
         BulkRequest.Builder br = new BulkRequest.Builder();
@@ -133,7 +144,11 @@ public class Elasticsearch8AsyncWriter<InputT> extends 
AsyncSinkWriter<InputT, O
             List<Operation> requestEntries,
             Consumer<List<Operation>> requestResult,
             Throwable error) {
-        LOG.debug("The BulkRequest of {} operation(s) has failed.", 
requestEntries.size());
+        LOG.warn(
+                "The BulkRequest of {} operation(s) has failed due to: {}",
+                requestEntries.size(),
+                error.getMessage());
+        LOG.debug("The BulkRequest has failed", error);
         numRecordsOutErrorsCounter.inc(requestEntries.size());
 
         if (isRetryable(error.getCause())) {
@@ -145,15 +160,17 @@ public class Elasticsearch8AsyncWriter<InputT> extends 
AsyncSinkWriter<InputT, O
             List<Operation> requestEntries,
             Consumer<List<Operation>> requestResult,
             BulkResponse response) {
+        LOG.debug("The BulkRequest has failed partially. Response: {}", 
response);
         ArrayList<Operation> failedItems = new ArrayList<>();
         for (int i = 0; i < response.items().size(); i++) {
             if (response.items().get(i).error() != null) {
-                numRecordsOutErrorsCounter.inc();
                 failedItems.add(requestEntries.get(i));
             }
         }
 
-        LOG.debug(
+        numRecordsOutErrorsCounter.inc(failedItems.size());
+        numRecordsSendPartialFailureCounter.inc(failedItems.size());
+        LOG.info(
                 "The BulkRequest with {} operation(s) has {} failure(s). It 
took {}ms",
                 requestEntries.size(),
                 failedItems.size(),

Reply via email to