23 changed files with 535 additions and 0 deletions
@ -0,0 +1,29 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.api; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
public interface MsqQueue { |
||||
|
|
||||
|
ListenableFuture<Void> put(TbMsg msg, UUID nodeId, long clusteredHash); |
||||
|
|
||||
|
ListenableFuture<Void> ack(TbMsg msg, UUID nodeId, long clusteredHash); |
||||
|
|
||||
|
Iterable<TbMsg> findUnprocessed(UUID nodeId, long clusteredHash); |
||||
|
} |
||||
@ -0,0 +1,29 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.queue.cassandra; |
||||
|
|
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.rule.engine.api.TbMsg; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
@Component |
||||
|
public class AckBuilder { |
||||
|
|
||||
|
public MsgAck build(TbMsg msg, UUID nodeId, long clusteredHash) { |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,82 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.queue.cassandra; |
||||
|
|
||||
|
import com.google.common.collect.Lists; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import org.springframework.beans.factory.annotation.Autowired; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.rule.engine.api.MsqQueue; |
||||
|
import org.thingsboard.rule.engine.api.TbMsg; |
||||
|
import org.thingsboard.rule.engine.queue.cassandra.repository.AckRepository; |
||||
|
import org.thingsboard.rule.engine.queue.cassandra.repository.MsgRepository; |
||||
|
import org.thingsboard.rule.engine.queue.cassandra.repository.ProcessedPartitionRepository; |
||||
|
|
||||
|
import java.util.Collections; |
||||
|
import java.util.List; |
||||
|
import java.util.Optional; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
@Component |
||||
|
public class CassandraMsqQueue implements MsqQueue { |
||||
|
|
||||
|
@Autowired |
||||
|
private MsgRepository msgRepository; |
||||
|
|
||||
|
@Autowired |
||||
|
private AckRepository ackRepository; |
||||
|
|
||||
|
@Autowired |
||||
|
private AckBuilder ackBuilder; |
||||
|
|
||||
|
@Autowired |
||||
|
private UnprocessedMsgFilter unprocessedMsgFilter; |
||||
|
|
||||
|
@Autowired |
||||
|
private ProcessedPartitionRepository processedPartitionRepository; |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Void> put(TbMsg msg, UUID nodeId, long clusteredHash) { |
||||
|
return msgRepository.save(msg, nodeId, clusteredHash, getPartition(msg)); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Void> ack(TbMsg msg, UUID nodeId, long clusteredHash) { |
||||
|
MsgAck ack = ackBuilder.build(msg, nodeId, clusteredHash); |
||||
|
return ackRepository.ack(ack); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public Iterable<TbMsg> findUnprocessed(UUID nodeId, long clusteredHash) { |
||||
|
List<TbMsg> unprocessedMsgs = Lists.newArrayList(); |
||||
|
for (Long partition : findUnprocessedPartitions(nodeId, clusteredHash)) { |
||||
|
Iterable<TbMsg> msgs = msgRepository.findMsgs(nodeId, clusteredHash, partition); |
||||
|
Iterable<MsgAck> acks = ackRepository.findAcks(nodeId, clusteredHash, partition); |
||||
|
unprocessedMsgs.addAll(unprocessedMsgFilter.filter(msgs, acks)); |
||||
|
} |
||||
|
return unprocessedMsgs; |
||||
|
} |
||||
|
|
||||
|
private List<Long> findUnprocessedPartitions(UUID nodeId, long clusteredHash) { |
||||
|
Optional<Long> lastPartition = processedPartitionRepository.findLastProcessedPartition(nodeId, clusteredHash); |
||||
|
return Collections.emptyList(); |
||||
|
} |
||||
|
|
||||
|
private long getPartition(TbMsg msg) { |
||||
|
return Long.MIN_VALUE; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,31 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.queue.cassandra; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
@Data |
||||
|
public class MsgAck { |
||||
|
|
||||
|
private final UUID msgId; |
||||
|
private final UUID nodeId; |
||||
|
private final long clusteredHash; |
||||
|
private final long partition; |
||||
|
private final long ts; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,28 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.queue.cassandra; |
||||
|
|
||||
|
import org.thingsboard.rule.engine.api.TbMsg; |
||||
|
|
||||
|
import java.util.Collection; |
||||
|
import java.util.Collections; |
||||
|
|
||||
|
public class UnprocessedMsgFilter { |
||||
|
|
||||
|
public Collection<TbMsg> filter(Iterable<TbMsg> msgs, Iterable<MsgAck> acks) { |
||||
|
return Collections.emptyList(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,28 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.queue.cassandra.repository; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import org.thingsboard.rule.engine.queue.cassandra.MsgAck; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
public interface AckRepository { |
||||
|
|
||||
|
ListenableFuture<Void> ack(MsgAck msgAck); |
||||
|
|
||||
|
Iterable<MsgAck> findAcks(UUID nodeId, long clusteredHash, long partition); |
||||
|
} |
||||
@ -0,0 +1,29 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.queue.cassandra.repository; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import org.thingsboard.rule.engine.api.TbMsg; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
public interface MsgRepository { |
||||
|
|
||||
|
ListenableFuture<Void> save(TbMsg msg, UUID nodeId, long clusteredHash, long partition); |
||||
|
|
||||
|
Iterable<TbMsg> findMsgs(UUID nodeId, long clusteredHash, long partition); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,12 @@ |
|||||
|
package org.thingsboard.rule.engine.queue.cassandra.repository; |
||||
|
|
||||
|
import java.util.Optional; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
public interface ProcessedPartitionRepository { |
||||
|
|
||||
|
void partitionProcessed(UUID nodeId, long clusteredHash, long partition); |
||||
|
|
||||
|
Optional<Long> findLastProcessedPartition(UUID nodeId, long clusteredHash); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,20 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.queue.jpa; |
||||
|
|
||||
|
//@todo-vp: implement
|
||||
|
public class SqlMsgQueue { |
||||
|
} |
||||
Loading…
Reference in new issue