From 537b010d474c489af5f63c1dd953f263bee8c5a4 Mon Sep 17 00:00:00 2001 From: penghui Date: Wed, 4 Jan 2023 20:35:36 +0800 Subject: [PATCH] [improve][broker] Do not print error log with stacktrace for 404 --- .../pulsar/broker/admin/AdminResource.java | 7 +++++++ .../admin/impl/SchemasResourceBase.java | 5 +++++ .../broker/admin/v1/SchemasResource.java | 20 +++++++++---------- .../broker/admin/v2/SchemasResource.java | 20 +++++++++---------- 4 files changed, 32 insertions(+), 20 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index 90eafa7c35ccc..09c9e2b032c6e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -846,6 +846,13 @@ protected static boolean isRedirectException(Throwable ex) { == Status.TEMPORARY_REDIRECT.getStatusCode(); } + protected static boolean isNotFoundException(Throwable ex) { + Throwable realCause = FutureUtil.unwrapCompletionException(ex); + return realCause instanceof WebApplicationException + && ((WebApplicationException) realCause).getResponse().getStatus() + == Status.NOT_FOUND.getStatusCode(); + } + protected static String getTopicNotFoundErrorMessage(String topic) { return String.format("Topic %s not found", topic); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java index f10bc913c2278..91a8140fde12b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java @@ -230,5 +230,10 @@ private CompletableFuture validateOwnershipAndOperationAsync(boolean autho .thenCompose(__ -> validateTopicOperationAsync(topicName, operation)); } + + protected boolean shouldPrintErrorLog(Throwable ex) { + return !isRedirectException(ex) && !isNotFoundException(ex); + } + private static final Logger log = LoggerFactory.getLogger(SchemasResourceBase.class); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/SchemasResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/SchemasResource.java index ca90eb4bac0c0..edc600707a120 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/SchemasResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/SchemasResource.java @@ -89,10 +89,10 @@ public void getSchema( ) { validateTopicName(tenant, cluster, namespace, topic); getSchemaAsync(authoritative) - .thenApply(schemaAndMetadata -> convertToSchemaResponse(schemaAndMetadata)) + .thenApply(this::convertToSchemaResponse) .thenApply(response::resume) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to get schema for topic {}", clientAppId(), topicName, ex); } resumeAsyncResponseExceptionally(response, ex); @@ -124,10 +124,10 @@ public void getSchema( ) { validateTopicName(tenant, cluster, namespace, topic); getSchemaAsync(authoritative, version) - .thenApply(schemaAndMetadata -> convertToSchemaResponse(schemaAndMetadata)) + .thenApply(this::convertToSchemaResponse) .thenAccept(response::resume) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to get schema for topic {} with version {}", clientAppId(), topicName, version, ex); } @@ -159,10 +159,10 @@ public void getAllSchemas( ) { validateTopicName(tenant, cluster, namespace, topic); getAllSchemasAsync(authoritative) - .thenApply(schemaAndMetadata -> convertToAllVersionsSchemaResponse(schemaAndMetadata)) + .thenApply(this::convertToAllVersionsSchemaResponse) .thenAccept(response::resume) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to get all schemas for topic {}", clientAppId(), topicName, ex); } resumeAsyncResponseExceptionally(response, ex); @@ -197,7 +197,7 @@ public void deleteSchema( response.resume(DeleteSchemaResponse.builder().version(getLongSchemaVersion(version)).build()); }) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to delete schemas for topic {}", clientAppId(), topicName, ex); } resumeAsyncResponseExceptionally(response, ex); @@ -252,7 +252,7 @@ public void postSchema( response.resume(Response.status(422, /* Unprocessable Entity */ root.getMessage()).build()); } else { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to post schemas for topic {}", clientAppId(), topicName, root); } resumeAsyncResponseExceptionally(response, ex); @@ -301,7 +301,7 @@ public void testCompatibility( .schemaCompatibilityStrategy(pair.getRight().name()).build()) .build())) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to test compatibility for topic {}", clientAppId(), topicName, ex); } resumeAsyncResponseExceptionally(response, ex); @@ -347,7 +347,7 @@ public void getVersionBySchema( getVersionBySchemaAsync(payload, authoritative) .thenAccept(version -> response.resume(LongSchemaVersionResponse.builder().version(version).build())) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to get version by schema for topic {}", clientAppId(), topicName, ex); } resumeAsyncResponseExceptionally(response, ex); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/SchemasResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/SchemasResource.java index 45d6bff42ca6a..952e75ebd392a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/SchemasResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/SchemasResource.java @@ -89,10 +89,10 @@ public void getSchema( @Suspended final AsyncResponse response) { validateTopicName(tenant, namespace, topic); getSchemaAsync(authoritative) - .thenApply(schemaAndMetadata -> convertToSchemaResponse(schemaAndMetadata)) + .thenApply(this::convertToSchemaResponse) .thenApply(response::resume) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to get schema for topic {}", clientAppId(), topicName, ex); } resumeAsyncResponseExceptionally(response, ex); @@ -122,10 +122,10 @@ public void getSchema( @Suspended final AsyncResponse response) { validateTopicName(tenant, namespace, topic); getSchemaAsync(authoritative, version) - .thenApply(schemaAndMetadata -> convertToSchemaResponse(schemaAndMetadata)) + .thenApply(this::convertToSchemaResponse) .thenAccept(response::resume) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to get schema for topic {} with version {}", clientAppId(), topicName, version, ex); } @@ -155,10 +155,10 @@ public void getAllSchemas( @Suspended final AsyncResponse response) { validateTopicName(tenant, namespace, topic); getAllSchemasAsync(authoritative) - .thenApply(schemaAndMetadata -> convertToAllVersionsSchemaResponse(schemaAndMetadata)) + .thenApply(this::convertToAllVersionsSchemaResponse) .thenAccept(response::resume) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to get all schemas for topic {}", clientAppId(), topicName, ex); } resumeAsyncResponseExceptionally(response, ex); @@ -191,7 +191,7 @@ public void deleteSchema( response.resume(DeleteSchemaResponse.builder().version(getLongSchemaVersion(version)).build()); }) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to delete schemas for topic {}", clientAppId(), topicName, ex); } resumeAsyncResponseExceptionally(response, ex); @@ -238,7 +238,7 @@ public void postSchema( response.resume(Response.status(422, /* Unprocessable Entity */ root.getMessage()).build()); } else { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to post schemas for topic {}", clientAppId(), topicName, root); } resumeAsyncResponseExceptionally(response, ex); @@ -278,7 +278,7 @@ public void testCompatibility( .schemaCompatibilityStrategy(pair.getRight().name()).build()) .build())) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to test compatibility for topic {}", clientAppId(), topicName, ex); } resumeAsyncResponseExceptionally(response, ex); @@ -315,7 +315,7 @@ public void getVersionBySchema( getVersionBySchemaAsync(payload, authoritative) .thenAccept(version -> response.resume(LongSchemaVersionResponse.builder().version(version).build())) .exceptionally(ex -> { - if (!isRedirectException(ex)) { + if (shouldPrintErrorLog(ex)) { log.error("[{}] Failed to get version by schema for topic {}", clientAppId(), topicName, ex); } resumeAsyncResponseExceptionally(response, ex);