Skip to content

Pubsub.Publisher: Failed to submit a listener notification task. Event loop shut down? #3326

Description

@mohamedazouz

I'm Receiving the following exception: Failed to submit a listener notification task. Event loop shut down?

class: io.grpc.netty.shaded.io.netty.util.concurrent.DefaultPromise	
exception:	
{	
         exception_class:	 java.lang.UnsupportedOperationException	
         stacktrace: 	java.lang.UnsupportedOperationException
         				at io.grpc.netty.shaded.io.netty.util.internal.shaded.org.jctools.queues.BaseMpscLinkedArrayQueue.iterator(BaseMpscLinkedArrayQueue.java:201)
						at java.util.AbstractCollection.remove(AbstractCollection.java:282)
						at io.grpc.netty.shaded.io.netty.util.concurrent.SingleThreadEventExecutor.removeTask(SingleThreadEventExecutor.java:339)
						at io.grpc.netty.shaded.io.netty.util.concurrent.SingleThreadEventExecutor.execute(SingleThreadEventExecutor.java:747)
						at io.grpc.netty.shaded.io.netty.util.concurrent.DefaultPromise.safeExecute(DefaultPromise.java:760)
						at io.grpc.netty.shaded.io.netty.util.concurrent.DefaultPromise.notifyListeners(DefaultPromise.java:428)
						at io.grpc.netty.shaded.io.netty.util.concurrent.DefaultPromise.cancel(DefaultPromise.java:315)
						at io.grpc.netty.shaded.io.netty.util.concurrent.ScheduledFutureTask.cancel(ScheduledFutureTask.java:151)
						at io.grpc.internal.ClientCallImpl.removeContextListenerAndCancelDeadlineFuture(ClientCallImpl.java:335)
						at io.grpc.internal.ClientCallImpl.access$1000(ClientCallImpl.java:63)
						at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl.close(ClientCallImpl.java:559)
						at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl.access$600(ClientCallImpl.java:478)
						at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1StreamClosed.runInContext(ClientCallImpl.java:590)
						at io.grpc.internal.ContextRunnable.run(ContextRunnable.java:37)
						at io.grpc.internal.SerializingExecutor.run(SerializingExecutor.java:123)
						at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
						at java.util.concurrent.FutureTask.run(FutureTask.java:266)
						at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180)
						at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293)
						at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
						at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
						at java.lang.Thread.run(Thread.java:748)	

 }	
file:	 DefaultPromise.java	
level:	 ERROR	
line_number:	 762	
logger_name:	 io.grpc.netty.shaded.io.netty.util.concurrent.DefaultPromise.rejectedExecution	
message:	 Failed to submit a listener notification task. Event loop shut down?	

here's how i configure publisher and how I publish:

final Publisher publisher = Publisher.newBuilder(topicName)
        .setChannelProvider(TopicAdminSettings.defaultTransportChannelProvider())
        .setCredentialsProvider(credentialsProvider) // loading cred. from file
        .build();

This is how i publish: I'm receiving batches contains maximum of 2 messages so publishing whenever new messages comes

try {

      for (final String message : messages) {
        final ByteString data = ByteString.copyFromUtf8(message);

        final PubsubMessage pubsubMessage = PubsubMessage.newBuilder()
            .setData(data)
            .build();
        final ApiFuture<String> messageIdFuture = publisher.publish(pubsubMessage);
        messageIdFutures.add(messageIdFuture);
      }

    } finally {
      final List<String> messageIds = ApiFutures.allAsList(messageIdFutures).get();

      for (final String messageId : messageIds) {
        LOG.info(MessageFormat.format("published with message ID: {0}", messageId));
      }
      if (publisher != null) {
        publisher.shutdown();
      }
      return messageIds;
    }

I don't know why I'm receiving above exception.

So I think it's related to concurrency issue I added the following to control concurrency but not sure if this will help me solving it or not.

final Publisher publisher = Publisher.newBuilder(topicName)
        .setChannelProvider(transportChannelProvider)
        .setCredentialsProvider(credentialsProvider)
        .setExecutorProvider(InstantiatingExecutorProvider.newBuilder()
        .setExecutorThreadCount(1).build())
        .build();

Activity

  1. added
    type: questionRequest for information or clarification. Not an issue.
    api: pubsubIssues related to the Pub/Sub API.
    priority: p2Moderately-important priority. Fix may not be included in next release.
    on Jun 1, 2018
  2. mohamedazouz commented on Jun 4, 2018

    @mohamedazouz
    Author

    With the same code: we're receiving com.google.api.gax.rpc.DeadlineExceededException: io.grpc.StatusRuntimeException: DEADLINE_EXCEEDED: deadline exceeded

  3. mohamedazouz commented on Jun 4, 2018

    @mohamedazouz
    Author

    I was using

    <dependency>
          <groupId>com.google.cloud</groupId>
          <artifactId>google-cloud-pubsub</artifactId>
          <version>0.47.0-beta</version>
        </dependency>
    

    Same issue happened with the newer version

    <dependency>
          <groupId>com.google.cloud</groupId>
          <artifactId>google-cloud-pubsub</artifactId>
          <version>1.31.0</version>
        </dependency>
    
  4. self-assigned this
    on Jun 4, 2018
  5. yihanzhen commented on Jun 5, 2018

    @yihanzhen
    Contributor

    Hi @mohamedazouz ,
    Sorry I was not able to repro this bug. Can you try it on v1.32.0 and let me know which line throws the exception? Also if it's possible can you try leaving everything for Publisher default (i.e. Publisher publisher = Publisher.newBuilder(topicName).build();) and see what happens?

    <dependency>
      <groupId>com.google.cloud</groupId>
      <artifactId>google-cloud-pubsub</artifactId>
      <version>1.32.0</version>
    </dependency>
    
  6. mohamedazouz commented on Jun 5, 2018

    @mohamedazouz
    Author

    @hzyi-google I will try the new version but I can't use the default publisher as i'm loading the cred. as well as the integration test, I will let you know if i got the same issue after upgrading version

  7. mohamedazouz commented on Jun 6, 2018

    @mohamedazouz
    Author

    @hzyi-google Here's the stack trace:

    Caused by: java.util.concurrent.ExecutionException: com.google.api.gax.rpc.DeadlineExceededException: io.grpc.StatusRuntimeException: DEADLINE_EXCEEDED: deadline exceeded after 8987261248ns
    at com.google.common.util.concurrent.AbstractFuture.getDoneValue(AbstractFuture.java:526)
    at com.google.common.util.concurrent.AbstractFuture.get(AbstractFuture.java:507)
    at com.google.common.util.concurrent.AbstractFuture$TrustedFuture.get(AbstractFuture.java:83)
    at com.google.common.util.concurrent.ForwardingFuture.get(ForwardingFuture.java:62)
    at  publisher.publish
    ... 1 more
    Caused by: com.google.api.gax.rpc.DeadlineExceededException: io.grpc.StatusRuntimeException: DEADLINE_EXCEEDED: deadline exceeded after 8987261248ns
    at com.google.api.gax.rpc.ApiExceptionFactory.createException(ApiExceptionFactory.java:51)
    at com.google.api.gax.grpc.GrpcApiExceptionFactory.create(GrpcApiExceptionFactory.java:72)
    at com.google.api.gax.grpc.GrpcApiExceptionFactory.create(GrpcApiExceptionFactory.java:60)
    at com.google.api.gax.grpc.GrpcExceptionCallable$ExceptionTransformingFuture.onFailure(GrpcExceptionCallable.java:95)
    at com.google.api.core.ApiFutures$1.onFailure(ApiFutures.java:61)
    at com.google.common.util.concurrent.Futures$CallbackListener.run(Futures.java:1341)
    at com.google.common.util.concurrent.MoreExecutors$DirectExecutor.execute(MoreExecutors.java:398)
    at com.google.common.util.concurrent.AbstractFuture.executeListener(AbstractFuture.java:1025)
    at com.google.common.util.concurrent.AbstractFuture.complete(AbstractFuture.java:866)
    at com.google.common.util.concurrent.AbstractFuture.setException(AbstractFuture.java:711)
    at io.grpc.stub.ClientCalls$GrpcFuture.setException(ClientCalls.java:492)
    at io.grpc.stub.ClientCalls$UnaryStreamToFuture.onClose(ClientCalls.java:467)
    at io.grpc.ForwardingClientCallListener.onClose(ForwardingClientCallListener.java:41)
    at io.grpc.internal.CensusStatsModule$StatsClientInterceptor$1$1.onClose(CensusStatsModule.java:684)
    at io.grpc.ForwardingClientCallListener.onClose(ForwardingClientCallListener.java:41)
    at io.grpc.internal.CensusTracingModule$TracingClientInterceptor$1$1.onClose(CensusTracingModule.java:391)
    at io.grpc.internal.ClientCallImpl.closeObserver(ClientCallImpl.java:475)
    at io.grpc.internal.ClientCallImpl.access$300(ClientCallImpl.java:63)
    at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl.close(ClientCallImpl.java:557)
    at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl.access$600(ClientCallImpl.java:478)
    at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1StreamClosed.runInContext(ClientCallImpl.java:590)
    at io.grpc.internal.ContextRunnable.run(ContextRunnable.java:37)
    at io.grpc.internal.SerializingExecutor.run(SerializingExecutor.java:123)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180)
    at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)
    Caused by: io.grpc.StatusRuntimeException: DEADLINE_EXCEEDED: deadline exceeded after 8987261248ns
    at io.grpc.Status.asRuntimeException(Status.java:526)
    ... 19 more
    
  8. pongad commented on Jun 6, 2018

    @pongad
    Contributor

    I think there are two different issues at play here. The UnsupportedOperationException seems like a gRPC bug. Could you let us know which version of gRPC you're using?

    The deadline exceeded happened because the RPC has been pending for ~10 seconds, much longer than we'd expect. How busy is the server? It's possible that the network is congested. Another possibility is that the executor that the Publisher is using is busy doing something else and cannot service the RPC in timely manner.

  9. mohamedazouz commented on Jun 7, 2018

    @mohamedazouz
    Author

    @pongad
    For the UnsupportedOperationException: I'm using google-cloud-pubsub with V: 1.32.0 which uses gRPC and it has the following dependency tree:

    [INFO] \- com.google.cloud:google-cloud-pubsub:jar:1.32.0:compile
    [INFO]    +- com.google.cloud:google-cloud-core:jar:1.32.0:compile
    [INFO]    |  +- com.google.http-client:google-http-client:jar:1.23.0:compile
    [INFO]    |  |  \- org.apache.httpcomponents:httpclient:jar:4.0.1:compile
    [INFO]    |  |     +- org.apache.httpcomponents:httpcore:jar:4.0.1:compile
    [INFO]    |  |     +- commons-logging:commons-logging:jar:1.1.1:compile
    [INFO]    |  |     \- commons-codec:commons-codec:jar:1.3:compile
    [INFO]    |  +- com.google.api:api-common:jar:1.6.0:compile
    [INFO]    |  +- com.google.api:gax:jar:1.27.0:compile
    [INFO]    |  |  +- org.threeten:threetenbp:jar:1.3.3:compile
    [INFO]    |  |  \- com.google.auth:google-auth-library-oauth2-http:jar:0.9.1:compile
    [INFO]    |  |     \- com.google.http-client:google-http-client-jackson2:jar:1.19.0:compile
    [INFO]    |  |        \- com.fasterxml.jackson.core:jackson-core:jar:2.1.3:compile
    [INFO]    |  +- com.google.protobuf:protobuf-java-util:jar:3.5.1:compile
    [INFO]    |  +- com.google.api.grpc:proto-google-common-protos:jar:1.11.0:compile
    [INFO]    |  \- com.google.api.grpc:proto-google-iam-v1:jar:0.12.0:compile
    [INFO]    +- com.google.cloud:google-cloud-core-grpc:jar:1.32.0:compile
    [INFO]    |  +- com.google.auth:google-auth-library-credentials:jar:0.9.1:compile
    [INFO]    |  +- com.google.protobuf:protobuf-java:jar:3.5.1:compile
    [INFO]    |  +- io.grpc:grpc-protobuf:jar:1.10.1:compile
    [INFO]    |  |  \- io.grpc:grpc-protobuf-lite:jar:1.10.1:compile
    [INFO]    |  +- io.grpc:grpc-context:jar:1.10.1:compile
    [INFO]    |  \- com.google.api:gax-grpc:jar:1.27.0:compile
    [INFO]    +- com.google.api.grpc:proto-google-cloud-pubsub-v1:jar:1.14.0:compile
    [INFO]    +- com.google.api.grpc:grpc-google-cloud-pubsub-v1:jar:1.14.0:compile
    [INFO]    +- io.grpc:grpc-netty-shaded:jar:1.10.1:compile
    [INFO]    |  \- io.grpc:grpc-core:jar:1.10.1:compile (version selected from constraint [1.10.1,1.10.1])
    [INFO]    |     +- io.opencensus:opencensus-api:jar:0.11.0:compile
    [INFO]    |     \- io.opencensus:opencensus-contrib-grpc-metrics:jar:0.11.0:compile
    [INFO]    +- io.grpc:grpc-stub:jar:1.10.1:compile
    [INFO]    +- io.grpc:grpc-auth:jar:1.10.1:compile
    [INFO]    \- com.google.auto.value:auto-value:jar:1.4:compile
    

    For the deadline exception: We're receiving at random times not only the peak times as We're receiving events to be pushed around ~2K/min with single message in each event. For the executor: I'm using the default Publisher settings besides for each event it creates a new publisher object (as publisher is triggered to publish messages comes on AWS Lambda invocation).

    I don't know if I have to set timeout as I can't change the behaviour of receiving single message, also what's the best practice here?

  10. mohamedazouz commented on Jun 7, 2018

    @mohamedazouz
    Author

    @hzyi-google I have tried the new version and i'm still receiving the same exceptions: #3326 (comment)

  11. pongad commented on Jun 7, 2018

    @pongad
    Contributor

    each event it creates a new publisher object

    Sorry I didn't realize this before. This might actually be the problem. The server imposes some limits on how many connections you can have open at once etc. It's possible you're running up against this limit: new messages can't be published because it's waiting for old connections to close -> it waits too long -> times out.

    Is it possible for events to share one (or a few) Publisher? The Publisher is written so that messages sent around the same time are batched and sent in one RPC, so sharing the publisher might decrease the number of connections you open significantly.

    The timeout might also be occurring because the network is congested, but I don't think 2K/min will do that.

  12. mohamedazouz commented on Jun 8, 2018

    @mohamedazouz
    Author

    @pongad Unfortunately, No i can't share one or few publisher.
    I think for this case i have to use http API instead or if there's a way to not use RPC ? What do you think?

  13. pongad commented on Jun 8, 2018

    @pongad
    Contributor

    Hmm. In that case I think it depends on why publisher can't be shared. If it's because you're publishing to many topics, maybe PublisherStub.publishCallable will help. You can create TopicAdminClient then call .getStub() to get to it. In this way we at least use the same gRPC connection to connect to the server.

    I'm not very familiar with Lambda, but if your environment doesn't let you keep a gRPC connection open, there might not be much of an advantage to use the gRPC client.

  14. mohamedazouz commented on Jun 8, 2018

    @mohamedazouz
    Author

    @pongad I will try to implement it using HTTP calls here, also i want to ask you if it's right thing to get the authentication heaer for each http requst from GoogleCredentials.getRequestMetadata() after getting Google as follows:

    GoogleCredentials.fromStream(new FileInputStream("cred_account_service.json"))
            .createScoped(Lists.newArrayList("https://www.googleapis.com/auth/cloud-platform","https://www.googleapis.com/auth/pubsub"));
    

    Please correct me what's the best way to fetch the token if this is not the correct way :)

  15. pongad commented on Jun 8, 2018

    @pongad
    Contributor

    That way of creating the credentials should be right. The getApplicationDefault method might be easier however.

  16. mohamedazouz commented on Jun 8, 2018

    @mohamedazouz
    Author

    @pongad Yes, I know but my question was this method GoogleCredentials.getRequestMetadata() after getting Google as follows: is the correct way to retrieve credentials Oath2 token?

  17. mohamedazouz commented on Jun 14, 2018

    @mohamedazouz
    Author

    After implementing custom http client, it worked fine.
    I will close the ticket and thanks for your support.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

api: pubsubIssues related to the Pub/Sub API.priority: p2Moderately-important priority. Fix may not be included in next release.type: questionRequest for information or clarification. Not an issue.

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions

    Sponsor
    SponsoredKunjungi sekarang
    Promo