From 948b74ad3697c85a91e044bbc2db7d042ee628ec Mon Sep 17 00:00:00 2001 From: Adam Date: Thu, 28 May 2026 09:36:30 +0200 Subject: [PATCH 01/12] TANGO-3103 : Permanent errors of bigquery-consumer mapped to 4xx http code --- .../consumer/sender/MessageSendingResult.java | 4 +++ .../sender/SingleMessageSendingResult.java | 6 +++++ .../GoogleBigQueryAppendCompleteCallback.java | 27 +++++++++++++++++-- 3 files changed, 35 insertions(+), 2 deletions(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/MessageSendingResult.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/MessageSendingResult.java index 45cf4a7496..1e3c3eedda 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/MessageSendingResult.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/MessageSendingResult.java @@ -36,6 +36,10 @@ static SingleMessageSendingResult failedResult(int statusCode) { return new SingleMessageSendingResult(statusCode); } + static SingleMessageSendingResult failedResult(int statusCode, Throwable cause) { + return new SingleMessageSendingResult(statusCode, cause); + } + static SingleMessageSendingResult ofStatusCode(int statusCode) { return new SingleMessageSendingResult(statusCode); } diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/SingleMessageSendingResult.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/SingleMessageSendingResult.java index 648ab2ed20..b83bbad3f0 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/SingleMessageSendingResult.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/SingleMessageSendingResult.java @@ -75,6 +75,12 @@ public String toString() { initializeForStatusCode(statusCode); } + SingleMessageSendingResult(int statusCode, Throwable failure) { + this.failure = failure; + initializeForStatusCode(statusCode); + this.loggable = true; + } + SingleMessageSendingResult(int statusCode, URI requestURI) { this(statusCode); this.requestUri = Optional.of(requestURI); diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java index fde0b39b10..5773ce2e14 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java @@ -3,12 +3,18 @@ import com.google.api.core.ApiFutureCallback; import com.google.cloud.bigquery.storage.v1.AppendRowsResponse; import com.google.cloud.bigquery.storage.v1.Exceptions; +import io.grpc.Status; import java.util.Objects; import java.util.concurrent.CompletableFuture; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import pl.allegro.tech.hermes.consumers.consumer.sender.MessageSendingResult; public class GoogleBigQueryAppendCompleteCallback implements ApiFutureCallback { + private static final Logger logger = + LoggerFactory.getLogger(GoogleBigQueryAppendCompleteCallback.class); + private final CompletableFuture resultFuture; public GoogleBigQueryAppendCompleteCallback( @@ -19,12 +25,29 @@ public GoogleBigQueryAppendCompleteCallback( @Override public void onFailure(Throwable t) { Exceptions.StorageException storageException = Exceptions.toStorageException(t); - resultFuture.complete( - MessageSendingResult.failedResult(Objects.requireNonNullElse(storageException, t))); + Throwable cause = Objects.requireNonNullElse(storageException, t); + + Integer httpStatusCode = mapToPermanentErrorHttpStatus(cause); + if (httpStatusCode != null) { + logger.warn("BigQuery permanent error mapped to HTTP {}: {}", httpStatusCode, cause.getMessage(), cause); + resultFuture.complete(MessageSendingResult.failedResult(httpStatusCode, cause)); + } else { + resultFuture.complete(MessageSendingResult.failedResult(cause)); + } } @Override public void onSuccess(AppendRowsResponse result) { resultFuture.complete(MessageSendingResult.succeededResult()); } + + private static Integer mapToPermanentErrorHttpStatus(Throwable cause) { + Status.Code grpcCode = Status.fromThrowable(cause).getCode(); + return switch (grpcCode) { + case NOT_FOUND -> 404; // Table does not exist + case PERMISSION_DENIED -> 403; // Technical user does not have permissions to write to the table + case INVALID_ARGUMENT -> 400; // Invalid message format i.e. microsecond timestamp value is sent to millisecond timestamp field + default -> null; + }; + } } From a0c38b7cc2581ec1332ba43660e55d7fcfe17e10 Mon Sep 17 00:00:00 2001 From: Adam Date: Thu, 28 May 2026 11:59:00 +0200 Subject: [PATCH 02/12] TANGO-3103 : Added logging of error and statusCode --- .../googlebigquery/GoogleBigQueryAppendCompleteCallback.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java index 5773ce2e14..dda9b13221 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java @@ -28,6 +28,7 @@ public void onFailure(Throwable t) { Throwable cause = Objects.requireNonNullElse(storageException, t); Integer httpStatusCode = mapToPermanentErrorHttpStatus(cause); + logger.info("BigQuery append failed with with code {} and error: {}",httpStatusCode, cause.getMessage(), cause); if (httpStatusCode != null) { logger.warn("BigQuery permanent error mapped to HTTP {}: {}", httpStatusCode, cause.getMessage(), cause); resultFuture.complete(MessageSendingResult.failedResult(httpStatusCode, cause)); @@ -43,11 +44,13 @@ public void onSuccess(AppendRowsResponse result) { private static Integer mapToPermanentErrorHttpStatus(Throwable cause) { Status.Code grpcCode = Status.fromThrowable(cause).getCode(); - return switch (grpcCode) { + Integer httpStatusCode = switch (grpcCode) { case NOT_FOUND -> 404; // Table does not exist case PERMISSION_DENIED -> 403; // Technical user does not have permissions to write to the table case INVALID_ARGUMENT -> 400; // Invalid message format i.e. microsecond timestamp value is sent to millisecond timestamp field default -> null; }; + logger.info("Mapping gRPC code {} to HTTP status code {}", grpcCode, httpStatusCode); + return httpStatusCode; } } From f41be80905065f4efc9ce87af19fb0f93c9e1b09 Mon Sep 17 00:00:00 2001 From: Adam Date: Thu, 28 May 2026 12:42:01 +0200 Subject: [PATCH 03/12] TANGO-3103 : Added custom http status code in error handling as well --- .../GoogleBigQueryAppendCompleteCallback.java | 2 +- .../sender/googlebigquery/GoogleBigQueryDataWriter.java | 7 ++++--- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java index dda9b13221..6f18768563 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java @@ -42,7 +42,7 @@ public void onSuccess(AppendRowsResponse result) { resultFuture.complete(MessageSendingResult.succeededResult()); } - private static Integer mapToPermanentErrorHttpStatus(Throwable cause) { + public static Integer mapToPermanentErrorHttpStatus(Throwable cause) { Status.Code grpcCode = Status.fromThrowable(cause).getCode(); Integer httpStatusCode = switch (grpcCode) { case NOT_FOUND -> 404; // Table does not exist diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java index c2fb5431d3..509d85dc3f 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java @@ -53,9 +53,9 @@ public void publish(T message, CompletableFuture resultFut .map(entry -> String.format("\t row %d: %s", entry.getKey(), entry.getValue())) .collect(Collectors.joining("\n")), e); - + Integer statusCode = GoogleBigQueryAppendCompleteCallback.mapToPermanentErrorHttpStatus(e); resultFuture.complete( - MessageSendingResult.failedResult(new GoogleBigQueryFailedAppendException(e))); + MessageSendingResult.failedResult(statusCode, new GoogleBigQueryFailedAppendException(e))); } catch (Exception e) { logger.warn( "Writer {} has failed to append rows to stream {} because of {}", @@ -63,7 +63,8 @@ public void publish(T message, CompletableFuture resultFut getStreamName(), e.getMessage(), e); - resultFuture.complete(MessageSendingResult.failedResult(e)); + Integer statusCode = GoogleBigQueryAppendCompleteCallback.mapToPermanentErrorHttpStatus(e); + resultFuture.complete(MessageSendingResult.failedResult(statusCode, e)); } } From 8aa085a8a764667b506b3549ed5d34b025052ff0 Mon Sep 17 00:00:00 2001 From: Adam Date: Thu, 11 Jun 2026 13:29:36 +0200 Subject: [PATCH 04/12] TANGO-3103 : Fixes after merge --- .../GoogleBigQueryDataWriter.java | 3 +- .../GoogleBigQueryFailedAppendException.java | 30 +++++++++++++++++++ 2 files changed, 31 insertions(+), 2 deletions(-) create mode 100644 hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java index 6a8741b36a..d13fc205b6 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java @@ -51,8 +51,7 @@ public void publish(T message, CompletableFuture resultFut .collect(Collectors.joining("\n")), e); Integer statusCode = GoogleBigQueryAppendCompleteCallback.mapToPermanentErrorHttpStatus(e); - throw e; - MessageSendingResult.failedResult(statusCode, new GoogleBigQueryFailedAppendException(e))); + MessageSendingResult.failedResult(statusCode, new GoogleBigQueryFailedAppendException(e)); } catch (Exception e) { logger.warn( "Writer {} has failed to append rows to stream {} because of {}", diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java new file mode 100644 index 0000000000..a4d7750ed5 --- /dev/null +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java @@ -0,0 +1,30 @@ +package pl.allegro.tech.hermes.consumers.consumer.sender.googlebigquery; + +import com.google.cloud.bigquery.storage.v1.Exceptions; +import java.util.Map; + +public class GoogleBigQueryFailedAppendException extends RuntimeException { + @Override + public String getMessage() { + Exceptions.AppendSerializtionError cause = + ((Exceptions.AppendSerializtionError) this.getCause()); + StringBuilder message = + new StringBuilder( + String.format( + "GoogleBigQuery Subscription has failed to append rows to stream %s", + cause.getStreamName())); + message.append(String.format("\n%s", super.getMessage())); + if (cause.getRowIndexToErrorMessage() != null) { + for (Map.Entry entry : cause.getRowIndexToErrorMessage().entrySet()) { + message.append( + String.format( + "\nGoogleBigQuery Subscription has failed because of %s", entry.getValue())); + } + } + return message.toString(); + } + + public GoogleBigQueryFailedAppendException(Exceptions.AppendSerializtionError cause) { + super(cause.getMessage(), cause); + } +} From 41b0a77fbafbb09dd8b640ba3f2bdc897727e6ed Mon Sep 17 00:00:00 2001 From: Adam Date: Fri, 12 Jun 2026 10:09:20 +0200 Subject: [PATCH 05/12] TANGO-3103 : Default 500 when any other error will appear --- .../googlebigquery/GoogleBigQueryAppendCompleteCallback.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java index 6f18768563..7bf670a981 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java @@ -48,7 +48,7 @@ public static Integer mapToPermanentErrorHttpStatus(Throwable cause) { case NOT_FOUND -> 404; // Table does not exist case PERMISSION_DENIED -> 403; // Technical user does not have permissions to write to the table case INVALID_ARGUMENT -> 400; // Invalid message format i.e. microsecond timestamp value is sent to millisecond timestamp field - default -> null; + default -> 500; }; logger.info("Mapping gRPC code {} to HTTP status code {}", grpcCode, httpStatusCode); return httpStatusCode; From 9e09a33fd8c7a968b013b9a987863b3afc1d1ae1 Mon Sep 17 00:00:00 2001 From: Adam Date: Fri, 12 Jun 2026 14:04:55 +0200 Subject: [PATCH 06/12] TANGO-3103 : Removed unnecessary loggings --- .../GoogleBigQueryAppendCompleteCallback.java | 22 ++++++------------- 1 file changed, 7 insertions(+), 15 deletions(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java index 7bf670a981..eae936f3ff 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java @@ -28,13 +28,7 @@ public void onFailure(Throwable t) { Throwable cause = Objects.requireNonNullElse(storageException, t); Integer httpStatusCode = mapToPermanentErrorHttpStatus(cause); - logger.info("BigQuery append failed with with code {} and error: {}",httpStatusCode, cause.getMessage(), cause); - if (httpStatusCode != null) { - logger.warn("BigQuery permanent error mapped to HTTP {}: {}", httpStatusCode, cause.getMessage(), cause); - resultFuture.complete(MessageSendingResult.failedResult(httpStatusCode, cause)); - } else { - resultFuture.complete(MessageSendingResult.failedResult(cause)); - } + resultFuture.complete(MessageSendingResult.failedResult(httpStatusCode, cause)); } @Override @@ -44,13 +38,11 @@ public void onSuccess(AppendRowsResponse result) { public static Integer mapToPermanentErrorHttpStatus(Throwable cause) { Status.Code grpcCode = Status.fromThrowable(cause).getCode(); - Integer httpStatusCode = switch (grpcCode) { - case NOT_FOUND -> 404; // Table does not exist - case PERMISSION_DENIED -> 403; // Technical user does not have permissions to write to the table - case INVALID_ARGUMENT -> 400; // Invalid message format i.e. microsecond timestamp value is sent to millisecond timestamp field - default -> 500; - }; - logger.info("Mapping gRPC code {} to HTTP status code {}", grpcCode, httpStatusCode); - return httpStatusCode; + return switch (grpcCode) { + case NOT_FOUND -> 404; // Table does not exist + case PERMISSION_DENIED -> 403; // Technical user does not have permissions to write to the table + case INVALID_ARGUMENT -> 400; // Invalid message format i.e. microsecond timestamp value is sent to millisecond timestamp field + default -> 500; + }; } } From 3feaf26d7fb139904fecf90c3ab30a190e8cbfa8 Mon Sep 17 00:00:00 2001 From: Adam Date: Fri, 12 Jun 2026 14:22:27 +0200 Subject: [PATCH 07/12] TANGO-3103 : Removed unused logger field --- .../googlebigquery/GoogleBigQueryAppendCompleteCallback.java | 3 --- 1 file changed, 3 deletions(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java index eae936f3ff..a918a23d3b 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java @@ -12,9 +12,6 @@ public class GoogleBigQueryAppendCompleteCallback implements ApiFutureCallback { - private static final Logger logger = - LoggerFactory.getLogger(GoogleBigQueryAppendCompleteCallback.class); - private final CompletableFuture resultFuture; public GoogleBigQueryAppendCompleteCallback( From 76cf7bc573ec8ced8100d7a48911c206eb561ee6 Mon Sep 17 00:00:00 2001 From: Adam Date: Fri, 12 Jun 2026 14:24:16 +0200 Subject: [PATCH 08/12] TANGO-3103 : Reformatted code --- .../GoogleBigQueryAppendCompleteCallback.java | 22 +++++++++---------- .../GoogleBigQueryFailedAppendException.java | 8 +++---- 2 files changed, 14 insertions(+), 16 deletions(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java index a918a23d3b..094ba56506 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java @@ -6,8 +6,6 @@ import io.grpc.Status; import java.util.Objects; import java.util.concurrent.CompletableFuture; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import pl.allegro.tech.hermes.consumers.consumer.sender.MessageSendingResult; public class GoogleBigQueryAppendCompleteCallback implements ApiFutureCallback { @@ -19,6 +17,16 @@ public GoogleBigQueryAppendCompleteCallback( this.resultFuture = resultFuture; } + public static Integer mapToPermanentErrorHttpStatus(Throwable cause) { + Status.Code grpcCode = Status.fromThrowable(cause).getCode(); + return switch (grpcCode) { + case NOT_FOUND -> 404; // Table does not exist + case PERMISSION_DENIED -> 403; // Technical user does not have permissions to write to the table + case INVALID_ARGUMENT -> 400; // Invalid message format i.e. microsecond timestamp value is sent to millisecond timestamp field + default -> 500; + }; + } + @Override public void onFailure(Throwable t) { Exceptions.StorageException storageException = Exceptions.toStorageException(t); @@ -32,14 +40,4 @@ public void onFailure(Throwable t) { public void onSuccess(AppendRowsResponse result) { resultFuture.complete(MessageSendingResult.succeededResult()); } - - public static Integer mapToPermanentErrorHttpStatus(Throwable cause) { - Status.Code grpcCode = Status.fromThrowable(cause).getCode(); - return switch (grpcCode) { - case NOT_FOUND -> 404; // Table does not exist - case PERMISSION_DENIED -> 403; // Technical user does not have permissions to write to the table - case INVALID_ARGUMENT -> 400; // Invalid message format i.e. microsecond timestamp value is sent to millisecond timestamp field - default -> 500; - }; - } } diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java index a4d7750ed5..8b83d83ea4 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java @@ -4,6 +4,10 @@ import java.util.Map; public class GoogleBigQueryFailedAppendException extends RuntimeException { + public GoogleBigQueryFailedAppendException(Exceptions.AppendSerializtionError cause) { + super(cause.getMessage(), cause); + } + @Override public String getMessage() { Exceptions.AppendSerializtionError cause = @@ -23,8 +27,4 @@ public String getMessage() { } return message.toString(); } - - public GoogleBigQueryFailedAppendException(Exceptions.AppendSerializtionError cause) { - super(cause.getMessage(), cause); - } } From 65c847362a0c06bbb82367b87c65159a4175a005 Mon Sep 17 00:00:00 2001 From: Adam Date: Fri, 12 Jun 2026 14:29:53 +0200 Subject: [PATCH 09/12] TANGO-3103 : Reformatted code v2 --- .../GoogleBigQueryAppendCompleteCallback.java | 13 +++--- .../GoogleBigQueryFailedAppendException.java | 42 +++++++++---------- 2 files changed, 29 insertions(+), 26 deletions(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java index 094ba56506..ef039c1a64 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java @@ -20,11 +20,14 @@ public GoogleBigQueryAppendCompleteCallback( public static Integer mapToPermanentErrorHttpStatus(Throwable cause) { Status.Code grpcCode = Status.fromThrowable(cause).getCode(); return switch (grpcCode) { - case NOT_FOUND -> 404; // Table does not exist - case PERMISSION_DENIED -> 403; // Technical user does not have permissions to write to the table - case INVALID_ARGUMENT -> 400; // Invalid message format i.e. microsecond timestamp value is sent to millisecond timestamp field - default -> 500; - }; + case NOT_FOUND -> 404; // Table does not exist + case PERMISSION_DENIED -> + 403; // Technical user does not have permissions to write to the table + case INVALID_ARGUMENT -> + 400; // Invalid message format i.e. microsecond timestamp value is sent to millisecond + // timestamp field + default -> 500; + }; } @Override diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java index 8b83d83ea4..d42e966302 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryFailedAppendException.java @@ -4,27 +4,27 @@ import java.util.Map; public class GoogleBigQueryFailedAppendException extends RuntimeException { - public GoogleBigQueryFailedAppendException(Exceptions.AppendSerializtionError cause) { - super(cause.getMessage(), cause); - } + public GoogleBigQueryFailedAppendException(Exceptions.AppendSerializtionError cause) { + super(cause.getMessage(), cause); + } - @Override - public String getMessage() { - Exceptions.AppendSerializtionError cause = - ((Exceptions.AppendSerializtionError) this.getCause()); - StringBuilder message = - new StringBuilder( - String.format( - "GoogleBigQuery Subscription has failed to append rows to stream %s", - cause.getStreamName())); - message.append(String.format("\n%s", super.getMessage())); - if (cause.getRowIndexToErrorMessage() != null) { - for (Map.Entry entry : cause.getRowIndexToErrorMessage().entrySet()) { - message.append( - String.format( - "\nGoogleBigQuery Subscription has failed because of %s", entry.getValue())); - } - } - return message.toString(); + @Override + public String getMessage() { + Exceptions.AppendSerializtionError cause = + ((Exceptions.AppendSerializtionError) this.getCause()); + StringBuilder message = + new StringBuilder( + String.format( + "GoogleBigQuery Subscription has failed to append rows to stream %s", + cause.getStreamName())); + message.append(String.format("\n%s", super.getMessage())); + if (cause.getRowIndexToErrorMessage() != null) { + for (Map.Entry entry : cause.getRowIndexToErrorMessage().entrySet()) { + message.append( + String.format( + "\nGoogleBigQuery Subscription has failed because of %s", entry.getValue())); + } } + return message.toString(); + } } From b5144ff06a740b48586e77e819ac89de3c83163d Mon Sep 17 00:00:00 2001 From: Adam Date: Fri, 12 Jun 2026 14:42:37 +0200 Subject: [PATCH 10/12] TANGO-3103 : Reformatted code v3 --- .../googlebigquery/GoogleBigQueryAppendCompleteCallback.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java index ef039c1a64..1cc2ce575c 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java @@ -25,7 +25,7 @@ public static Integer mapToPermanentErrorHttpStatus(Throwable cause) { 403; // Technical user does not have permissions to write to the table case INVALID_ARGUMENT -> 400; // Invalid message format i.e. microsecond timestamp value is sent to millisecond - // timestamp field + // timestamp field default -> 500; }; } From a1a359fa0d8f8cbe90bbfbeb0318a0bbc519e103 Mon Sep 17 00:00:00 2001 From: Adam Szorcz <87751153+adamallegro@users.noreply.github.com> Date: Tue, 30 Jun 2026 08:13:29 +0200 Subject: [PATCH 11/12] TANGO-3103: Potential fix for pull request finding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- .../sender/googlebigquery/GoogleBigQueryDataWriter.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java index d13fc205b6..c5f0fe3d60 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java @@ -50,8 +50,9 @@ public void publish(T message, CompletableFuture resultFut .map(entry -> String.format("\t row %d: %s", entry.getKey(), entry.getValue())) .collect(Collectors.joining("\n")), e); - Integer statusCode = GoogleBigQueryAppendCompleteCallback.mapToPermanentErrorHttpStatus(e); - MessageSendingResult.failedResult(statusCode, new GoogleBigQueryFailedAppendException(e)); + int statusCode = GoogleBigQueryAppendCompleteCallback.mapToPermanentErrorHttpStatus(e); + resultFuture.complete( + MessageSendingResult.failedResult(statusCode, new GoogleBigQueryFailedAppendException(e))); } catch (Exception e) { logger.warn( "Writer {} has failed to append rows to stream {} because of {}", From 0362ed4d7e58a8a4b8787f33980e1ef7dde37aaf Mon Sep 17 00:00:00 2001 From: Adam Date: Thu, 13 Aug 2026 07:47:01 +0200 Subject: [PATCH 12/12] TANGO-3103 : Applied review changes --- .../googlebigquery/GoogleBigQueryAppendCompleteCallback.java | 4 ++-- .../sender/googlebigquery/GoogleBigQueryDataWriter.java | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java index 1cc2ce575c..8b8cdfedd6 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryAppendCompleteCallback.java @@ -17,7 +17,7 @@ public GoogleBigQueryAppendCompleteCallback( this.resultFuture = resultFuture; } - public static Integer mapToPermanentErrorHttpStatus(Throwable cause) { + public static int mapToErrorHttpStatus(Throwable cause) { Status.Code grpcCode = Status.fromThrowable(cause).getCode(); return switch (grpcCode) { case NOT_FOUND -> 404; // Table does not exist @@ -35,7 +35,7 @@ public void onFailure(Throwable t) { Exceptions.StorageException storageException = Exceptions.toStorageException(t); Throwable cause = Objects.requireNonNullElse(storageException, t); - Integer httpStatusCode = mapToPermanentErrorHttpStatus(cause); + Integer httpStatusCode = mapToErrorHttpStatus(cause); resultFuture.complete(MessageSendingResult.failedResult(httpStatusCode, cause)); } diff --git a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java index c5f0fe3d60..54a0f88687 100644 --- a/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java +++ b/hermes-consumers/src/main/java/pl/allegro/tech/hermes/consumers/consumer/sender/googlebigquery/GoogleBigQueryDataWriter.java @@ -50,7 +50,7 @@ public void publish(T message, CompletableFuture resultFut .map(entry -> String.format("\t row %d: %s", entry.getKey(), entry.getValue())) .collect(Collectors.joining("\n")), e); - int statusCode = GoogleBigQueryAppendCompleteCallback.mapToPermanentErrorHttpStatus(e); + int statusCode = GoogleBigQueryAppendCompleteCallback.mapToErrorHttpStatus(e); resultFuture.complete( MessageSendingResult.failedResult(statusCode, new GoogleBigQueryFailedAppendException(e))); } catch (Exception e) { @@ -60,7 +60,7 @@ public void publish(T message, CompletableFuture resultFut getStreamName(), e.getMessage(), e); - Integer statusCode = GoogleBigQueryAppendCompleteCallback.mapToPermanentErrorHttpStatus(e); + Integer statusCode = GoogleBigQueryAppendCompleteCallback.mapToErrorHttpStatus(e); resultFuture.complete(MessageSendingResult.failedResult(statusCode, e)); } }