This is an automated email from the ASF dual-hosted git repository.
duhengforever pushed a change to branch mqtt
in repository https://gitbox.apache.org/repos/asf/rocketmq.git.
from 4f8cd91 Merge pull request #1197 from xiangwangcheng/mqtt
new e010ca7 polish data structure of IOTClientManager;add message forward
logic
new f8da6ff optimize Client of IOT; add in-flight window
new 25ba6fa optimize logic of push message to consumers;add
MQTTSession(extends Client)
new 45d8d04 add logic of persist consumeOffset;polish logic of PUBLISH
messages to clients(consumers)
new b1a0b48 Merge pull request #1209 from xiangwangcheng/mqtt
The 934 revisions listed above as "new" are entirely new to this
repository and will be described in separate emails. The revisions
listed as "add" were already present in the repository and have only
been added to this reference.
Summary of changes:
.../client/consumer/DefaultMQPullConsumerTest.java | 2 +-
.../client/consumer/DefaultMQPushConsumerTest.java | 2 +-
common/pom.xml | 4 +-
.../org/apache/rocketmq/common/MqttConfig.java | 11 +-
.../org/apache/rocketmq/common/client/Client.java | 48 +-
mqtt/pom.xml | 13 +
.../rocketmq/mqtt/client/IOTClientManagerImpl.java | 63 +-
.../rocketmq/mqtt/client/InFlightMessage.java | 62 ++
.../apache/rocketmq/mqtt/client/MQTTSession.java | 255 ++++++++
.../rocketmq/mqtt/constant/MqttConstant.java | 1 +
.../rocketmq/mqtt/mqtthandler/MessageHandler.java | 58 +-
.../impl/MqttConnectMessageHandler.java | 14 +-
.../impl/MqttDisconnectMessageHandler.java | 4 +-
.../impl/MqttMessageForwardHandler.java | 98 +++
.../mqtthandler/impl/MqttMessageForwarder.java | 46 --
.../impl/MqttPingreqMessageHandler.java | 7 +-
.../mqtthandler/impl/MqttPubackMessageHandler.java | 28 +
.../impl/MqttPublishMessageHandler.java | 199 +++---
.../impl/MqttSubscribeMessageHandler.java | 35 +-
.../impl/MqttUnsubscribeMessagHandler.java | 43 +-
.../processor/DefaultMqttMessageProcessor.java | 29 +-
.../mqtt/processor/InnerMqttMessageProcessor.java | 137 +++++
.../mqtt/service/impl/MqttPushServiceImpl.java | 164 -----
.../service/impl/MqttScheduledServiceImpl.java | 78 +++
.../apache/rocketmq/mqtt/task/MqttPushTask.java | 171 ++++++
.../rocketmq/mqtt/transfer/TransferDataQos1.java | 21 +-
.../orderedexecutor/BoundedExecutorService.java | 108 ++++
.../mqtt/util/orderedexecutor/Counter.java | 50 ++
.../rocketmq/mqtt/util/orderedexecutor/Gauge.java | 28 +
.../mqtt/util/orderedexecutor/MathUtils.java | 90 +++
.../mqtt/util/orderedexecutor/MdcUtils.java | 39 ++
.../mqtt/util/orderedexecutor/NullStatsLogger.java | 129 ++++
.../mqtt/util/orderedexecutor/OpStatsData.java | 76 +++
.../mqtt/util/orderedexecutor/OpStatsLogger.java | 63 ++
.../mqtt/util/orderedexecutor/OrderedExecutor.java | 684 +++++++++++++++++++++
.../mqtt/util/orderedexecutor/SafeRunnable.java | 95 +++
.../mqtt/util/orderedexecutor/StatsLogger.java | 76 +++
.../remoting/transport/mqtt/MqttHeader.java | 4 +-
.../org/apache/rocketmq/snode/SnodeController.java | 34 +-
39 files changed, 2632 insertions(+), 437 deletions(-)
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/client/InFlightMessage.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/client/MQTTSession.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/mqtthandler/impl/MqttMessageForwardHandler.java
delete mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/mqtthandler/impl/MqttMessageForwarder.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/processor/InnerMqttMessageProcessor.java
delete mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/service/impl/MqttPushServiceImpl.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/service/impl/MqttScheduledServiceImpl.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/task/MqttPushTask.java
copy
common/src/main/java/org/apache/rocketmq/common/protocol/header/GetTopicStatsInfoRequestHeader.java
=> mqtt/src/main/java/org/apache/rocketmq/mqtt/transfer/TransferDataQos1.java
(69%)
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/util/orderedexecutor/BoundedExecutorService.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/util/orderedexecutor/Counter.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/util/orderedexecutor/Gauge.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/util/orderedexecutor/MathUtils.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/util/orderedexecutor/MdcUtils.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/util/orderedexecutor/NullStatsLogger.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/util/orderedexecutor/OpStatsData.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/util/orderedexecutor/OpStatsLogger.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/util/orderedexecutor/OrderedExecutor.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/util/orderedexecutor/SafeRunnable.java
create mode 100644
mqtt/src/main/java/org/apache/rocketmq/mqtt/util/orderedexecutor/StatsLogger.java