0

使用 Azure 服务消息总线主题和订阅。

创建了一个主题aaaa和订阅zzzsubscription- 会话感知。

使用以下代码向具有会话 ID 的主题发送消息:

String senderString="Endpoint=sb://xxxx.servicebus.windows.net/;SharedAccessKeyName=sampler-sender-only-policy;SharedAccessKey=mmmm;EntityPath=aaaa";
ConnectionStringBuilder connectionStringBuilder = new ConnectionStringBuilder(senderString,
                "aaaa");
TopicClient client=new TopicClient(connectionStringBuilder);
for(int session =1 ;session<=10;session++){
        for(int message =1 ;message<=10;message++){
                Message sendMessage=new Message("message "+message);
                sendMessage.setMessageId(UUID.randomUUID().toString());
                sendMessage.setSessionId("Session "+session );
                client.sendAsync(sendMessage);
        }
}

使用以下代码从订阅中读取消息:

String listenerString = "Endpoint=sb://xxxx.servicebus.windows.net/;SharedAccessKeyName=sample-listen-only-policy;SharedAccessKey=yyyy;EntityPath=aaaa";
ConnectionStringBuilder connectionStringBuilderListen = new ConnectionStringBuilder(listenerString,
                "aaaa/subscriptions/zzzsubscription");
SubscriptionClient subscriptionClient = new SubscriptionClient(connectionStringBuilderListen, ReceiveMode.PEEKLOCK);

subscriptionClient.registerSessionHandler(new ISessionHandler() {
        @Override
        public CompletableFuture<Void> onMessageAsync(IMessageSession session, IMessage message) {
                System.out.println(message.getSessionId()+" - "+ new String(message.getBody()));
                return subscriptionClient.completeAsync(message.getLockToken());
        }
        @Override
         public void notifyException(Throwable throwable, ExceptionPhase exceptionPhase) {
                System.out.printf(exceptionPhase + "-" + throwable.getMessage());
            }
        @Override
        public CompletableFuture<Void> OnCloseSessionAsync(IMessageSession session) {
                return subscriptionClient.closeAsync();
        }
});

在该行获取以下异常 return subscriptionClient.completeAsync(message.getLockToken());

2018-02-23 11:05:15,471 [onPool-worker-2] - ERROR MessageAndSessionPump          - onMessage with message containing sequence number '15481123719086103' threw exception
java.lang.UnsupportedOperationException: Receiver not created. Registering a MessageHandler creates a receiver.
    at com.microsoft.azure.servicebus.MessageAndSessionPump.checkInnerReceiveCreated(MessageAndSessionPump.java:712)
    at com.microsoft.azure.servicebus.MessageAndSessionPump.completeAsync(MessageAndSessionPump.java:636)
    at com.microsoft.azure.servicebus.SubscriptionClient.completeAsync(SubscriptionClient.java:208)
    at com.microsoft.azure.servicebus.samples.topicsgettingstarted.TopicsGettingStarted$1.onMessageAsync(TopicsGettingStarted.java:31)

使用以下版本的 Azure 服务总线 SDK:

<dependency>
    <groupId>com.microsoft.azure</groupId>
    <artifactId>azure-servicebus</artifactId>
    <version>1.1.1</version>
</dependency>

编辑1:

如果我添加以下内容,则到上面的代码:

SessionHandlerOptions sessionHandlerOptions = new SessionHandlerOptions(1, true, Duration.ofMinutes(5));
subscriptionClient.registerSessionHandler(new ISessionHandler() {...}, sessionHandlerOptions);

并将 onMessageAsync 中的返回更改为以下内容:

return CompletableFuture.completedFuture(null);

然后我得到以下错误。收到了一些消息,但我收到以下错误并且程序挂起:

2018-02-23 11:37:17,814 [24-7360c77df5c2] - ERROR RequestResponseLink            - Opening internal send link of requestresponselink to sample-order-status/subscriptions/sample-hybris-order-status-subscription/$management failed.
com.microsoft.azure.servicebus.primitives.AuthorizationFailedException: Unauthorized access. 'Send' claim(s) are required to perform this operation. Resource: 'sb://ashokgoli.servicebus.windows.net/sample-order-status/subscriptions/sample-hybris-order-status-subscription/$management'. TrackingId:cbe7e9aa831a455abdaa17aacb814497_G27, SystemTracker:gateway6, Timestamp:2/23/2018 5:37:17 PM
    at com.microsoft.azure.servicebus.primitives.ExceptionUtil.toException(ExceptionUtil.java:50)
    at com.microsoft.azure.servicebus.primitives.RequestResponseLink$InternalSender.onClose(RequestResponseLink.java:743)
    at com.microsoft.azure.servicebus.amqp.BaseLinkHandler.processOnClose(BaseLinkHandler.java:68)
    at com.microsoft.azure.servicebus.amqp.BaseLinkHandler.onLinkRemoteClose(BaseLinkHandler.java:42)
    at org.apache.qpid.proton.engine.BaseHandler.handle(BaseHandler.java:176)
    at org.apache.qpid.proton.engine.impl.EventImpl.dispatch(EventImpl.java:108)
    at org.apache.qpid.proton.reactor.impl.ReactorImpl.dispatch(ReactorImpl.java:309)
    at org.apache.qpid.proton.reactor.impl.ReactorImpl.process(ReactorImpl.java:276)
    at com.microsoft.azure.servicebus.primitives.MessagingFactory$RunReactor.run(MessagingFactory.java:451)
    at java.base/java.lang.Thread.run(Thread.java:844)

通过将发送声明添加到只听策略解决了上述问题:https ://github.com/Azure/azure-service-bus-java/issues/110

2018-02-23 11:37:47,541 [nPool-worker-15] - ERROR RequestResponseLinkcache       - Creating requestresponselink to 'sample-order-status/subscriptions/sample-hybris-order-status-subscription/$management' failed.
com.microsoft.azure.servicebus.primitives.TimeoutException: Open operation on RequestResponseLink(05a91c-RequestResponse) on Entity(sample-order-status/subscriptions/sample-hybris-order-status-subscription/$management) timed out at 2018-02-23T11:37:47.540528-06:00[America/Chicago].
    at com.microsoft.azure.servicebus.primitives.RequestResponseLink$1.run(RequestResponseLink.java:77)
    at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:514)
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
    at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:299)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1167)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:641)
    at java.base/java.lang.Thread.run(Thread.java:844)
4

1 回答 1

0

在 OnCloseSessionAsync 上返回未来。您可以参考以下链接以供参考。

https://github.com/Azure/azure-service-bus-java/issues/187

于 2018-03-02T10:41:28.677 回答