()
| 113 | } |
| 114 | |
| 115 | private void tryPushMessage() { |
| 116 | if (null == pendingPushMessage && !messageQueue.isEmpty() && bindSocket != null) { |
| 117 | pendingPushMessage = new PushMessage(); |
| 118 | pendingPushMessage.Argument.setTopic(topic); |
| 119 | pendingPushMessage.Argument.setPartitionIndex(partitionIndex); |
| 120 | pendingPushMessage.Argument.setSessionId(bindSessionId); |
| 121 | var message = messageQueue.peek(); |
| 122 | pendingPushMessage.Argument.setMessage(message); |
| 123 | pendingPushMessage.Send(bindSocket, (p) -> { |
| 124 | lock(); |
| 125 | try { |
| 126 | loadCounter.incrementAndGet(); // 处理失败也进行计数。 |
| 127 | |
| 128 | if (pendingPushMessage.getResultCode() == 0) { |
| 129 | messageQueue.poll(); |
| 130 | fileWithIndex.increaseFirstMessageId(); |
| 131 | tryStartBackgroundFill(); |
| 132 | } |
| 133 | |
| 134 | // 不管是否失败,都尝试重新pushMessage。出错的时候要不要随机延迟一下再重试? |
| 135 | pendingPushMessage = null; |
| 136 | tryPushMessage(); |
| 137 | } finally { |
| 138 | unlock(); |
| 139 | } |
| 140 | return 0; |
| 141 | }); |
| 142 | } |
| 143 | } |
| 144 | |
| 145 | public void bind(long sessionId, AsyncSocket socket) { |
| 146 | lock(); |
no test coverage detected