52 changed files with 454 additions and 622 deletions
@ -0,0 +1,80 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2018 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.server.common.msg; |
||||
|
|
||||
|
import com.google.protobuf.ByteString; |
||||
|
import com.google.protobuf.InvalidProtocolBufferException; |
||||
|
import lombok.Data; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.EntityIdFactory; |
||||
|
import org.thingsboard.server.common.msg.gen.MsgProtos; |
||||
|
|
||||
|
import java.io.Serializable; |
||||
|
import java.nio.ByteBuffer; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 13.01.18. |
||||
|
*/ |
||||
|
@Data |
||||
|
public final class TbMsg implements Serializable { |
||||
|
|
||||
|
private final UUID id; |
||||
|
private final String type; |
||||
|
private final EntityId originator; |
||||
|
private final TbMsgMetaData metaData; |
||||
|
|
||||
|
private final byte[] data; |
||||
|
|
||||
|
public static ByteBuffer toBytes(TbMsg msg) { |
||||
|
MsgProtos.TbMsgProto.Builder builder = MsgProtos.TbMsgProto.newBuilder(); |
||||
|
builder.setId(msg.getId().toString()); |
||||
|
builder.setType(msg.getType()); |
||||
|
if (msg.getOriginator() != null) { |
||||
|
builder.setEntityType(msg.getOriginator().getEntityType().name()); |
||||
|
builder.setEntityId(msg.getOriginator().getId().toString()); |
||||
|
} |
||||
|
|
||||
|
if (msg.getMetaData() != null) { |
||||
|
MsgProtos.TbMsgProto.TbMsgMetaDataProto.Builder metadataBuilder = MsgProtos.TbMsgProto.TbMsgMetaDataProto.newBuilder(); |
||||
|
metadataBuilder.putAllData(msg.getMetaData().getData()); |
||||
|
builder.addMetaData(metadataBuilder.build()); |
||||
|
} |
||||
|
|
||||
|
builder.setData(ByteString.copyFrom(msg.getData())); |
||||
|
byte[] bytes = builder.build().toByteArray(); |
||||
|
return ByteBuffer.wrap(bytes); |
||||
|
} |
||||
|
|
||||
|
public static TbMsg fromBytes(ByteBuffer buffer) { |
||||
|
try { |
||||
|
MsgProtos.TbMsgProto proto = MsgProtos.TbMsgProto.parseFrom(buffer.array()); |
||||
|
TbMsgMetaData metaData = new TbMsgMetaData(); |
||||
|
if (proto.getMetaDataCount() > 0) { |
||||
|
metaData.setData(proto.getMetaData(0).getDataMap()); |
||||
|
} |
||||
|
|
||||
|
EntityId entityId = null; |
||||
|
if (proto.getEntityId() != null) { |
||||
|
entityId = EntityIdFactory.getByTypeAndId(proto.getEntityType(), proto.getEntityId()); |
||||
|
} |
||||
|
|
||||
|
return new TbMsg(UUID.fromString(proto.getId()), proto.getType(), entityId, metaData, proto.getData().toByteArray()); |
||||
|
} catch (InvalidProtocolBufferException e) { |
||||
|
throw new IllegalStateException("Could not parse protobuf for TbMsg", e); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -1,62 +1,62 @@ |
|||||
/** |
/** |
||||
* Copyright © 2016-2017 The Thingsboard Authors |
* Copyright © 2016-2018 The Thingsboard Authors |
||||
* <p> |
* |
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
* you may not use this file except in compliance with the License. |
* you may not use this file except in compliance with the License. |
||||
* You may obtain a copy of the License at |
* You may obtain a copy of the License at |
||||
* <p> |
* |
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
* <p> |
* |
||||
* Unless required by applicable law or agreed to in writing, software |
* Unless required by applicable law or agreed to in writing, software |
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
* See the License for the specific language governing permissions and |
* See the License for the specific language governing permissions and |
||||
* limitations under the License. |
* limitations under the License. |
||||
*/ |
*/ |
||||
package org.thingsboard.rule.engine.queue.cassandra.repository.impl; |
package org.thingsboard.server.dao.service.queue.cassandra.repository.impl; |
||||
|
|
||||
import com.datastax.driver.core.*; |
import com.datastax.driver.core.*; |
||||
import com.google.common.base.Function; |
import com.google.common.base.Function; |
||||
import com.google.common.util.concurrent.Futures; |
import com.google.common.util.concurrent.Futures; |
||||
import com.google.common.util.concurrent.ListenableFuture; |
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
import org.springframework.stereotype.Component; |
import org.springframework.stereotype.Component; |
||||
import org.thingsboard.rule.engine.queue.cassandra.MsgAck; |
import org.thingsboard.server.dao.nosql.CassandraAbstractDao; |
||||
import org.thingsboard.rule.engine.queue.cassandra.repository.AckRepository; |
import org.thingsboard.server.dao.service.queue.cassandra.MsgAck; |
||||
|
import org.thingsboard.server.dao.service.queue.cassandra.repository.AckRepository; |
||||
|
import org.thingsboard.server.dao.util.NoSqlDao; |
||||
|
|
||||
import java.util.ArrayList; |
import java.util.ArrayList; |
||||
import java.util.List; |
import java.util.List; |
||||
import java.util.UUID; |
import java.util.UUID; |
||||
|
|
||||
@Component |
@Component |
||||
public class CassandraAckRepository extends SimpleAbstractCassandraDao implements AckRepository { |
@NoSqlDao |
||||
|
public class CassandraAckRepository extends CassandraAbstractDao implements AckRepository { |
||||
|
|
||||
private final int ackQueueTtl; |
@Value("${cassandra.queue.ack.ttl}") |
||||
|
private int ackQueueTtl; |
||||
public CassandraAckRepository(Session session, int ackQueueTtl) { |
|
||||
super(session); |
|
||||
this.ackQueueTtl = ackQueueTtl; |
|
||||
} |
|
||||
|
|
||||
@Override |
@Override |
||||
public ListenableFuture<Void> ack(MsgAck msgAck) { |
public ListenableFuture<Void> ack(MsgAck msgAck) { |
||||
String insert = "INSERT INTO msg_ack_queue (node_id, clustered_hash, partition, msg_id) VALUES (?, ?, ?, ?) USING TTL ?"; |
String insert = "INSERT INTO msg_ack_queue (node_id, cluster_partition, ts_partition, msg_id) VALUES (?, ?, ?, ?) USING TTL ?"; |
||||
PreparedStatement statement = prepare(insert); |
PreparedStatement statement = prepare(insert); |
||||
BoundStatement boundStatement = statement.bind(msgAck.getNodeId(), msgAck.getClusteredHash(), |
BoundStatement boundStatement = statement.bind(msgAck.getNodeId(), msgAck.getClusteredPartition(), |
||||
msgAck.getPartition(), msgAck.getMsgId(), ackQueueTtl); |
msgAck.getTsPartition(), msgAck.getMsgId(), ackQueueTtl); |
||||
ResultSetFuture resultSetFuture = executeAsyncWrite(boundStatement); |
ResultSetFuture resultSetFuture = executeAsyncWrite(boundStatement); |
||||
return Futures.transform(resultSetFuture, (Function<ResultSet, Void>) input -> null); |
return Futures.transform(resultSetFuture, (Function<ResultSet, Void>) input -> null); |
||||
} |
} |
||||
|
|
||||
@Override |
@Override |
||||
public List<MsgAck> findAcks(UUID nodeId, long clusteredHash, long partition) { |
public List<MsgAck> findAcks(UUID nodeId, long clusterPartition, long tsPartition) { |
||||
String select = "SELECT msg_id FROM msg_ack_queue WHERE " + |
String select = "SELECT msg_id FROM msg_ack_queue WHERE " + |
||||
"node_id = ? AND clustered_hash = ? AND partition = ?"; |
"node_id = ? AND cluster_partition = ? AND ts_partition = ?"; |
||||
PreparedStatement statement = prepare(select); |
PreparedStatement statement = prepare(select); |
||||
BoundStatement boundStatement = statement.bind(nodeId, clusteredHash, partition); |
BoundStatement boundStatement = statement.bind(nodeId, clusterPartition, tsPartition); |
||||
ResultSet rows = executeRead(boundStatement); |
ResultSet rows = executeRead(boundStatement); |
||||
List<MsgAck> msgs = new ArrayList<>(); |
List<MsgAck> msgs = new ArrayList<>(); |
||||
for (Row row : rows) { |
for (Row row : rows) { |
||||
msgs.add(new MsgAck(row.getUUID("msg_id"), nodeId, clusteredHash, partition)); |
msgs.add(new MsgAck(row.getUUID("msg_id"), nodeId, clusterPartition, tsPartition)); |
||||
} |
} |
||||
return msgs; |
return msgs; |
||||
} |
} |
||||
@ -0,0 +1,63 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2018 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.server.dao.service.queue.cassandra.repository.impl; |
||||
|
|
||||
|
import com.datastax.driver.core.*; |
||||
|
import com.google.common.base.Function; |
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
import org.thingsboard.server.dao.nosql.CassandraAbstractDao; |
||||
|
import org.thingsboard.server.dao.service.queue.cassandra.repository.MsgRepository; |
||||
|
import org.thingsboard.server.dao.util.NoSqlDao; |
||||
|
|
||||
|
import java.util.ArrayList; |
||||
|
import java.util.List; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
@Component |
||||
|
@NoSqlDao |
||||
|
public class CassandraMsgRepository extends CassandraAbstractDao implements MsgRepository { |
||||
|
|
||||
|
@Value("${cassandra.queue.msg.ttl}") |
||||
|
private int msqQueueTtl; |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Void> save(TbMsg msg, UUID nodeId, long clusterPartition, long tsPartition, long msgTs) { |
||||
|
String insert = "INSERT INTO msg_queue (node_id, cluster_partition, ts_partition, ts, msg) VALUES (?, ?, ?, ?, ?) USING TTL ?"; |
||||
|
PreparedStatement statement = prepare(insert); |
||||
|
BoundStatement boundStatement = statement.bind(nodeId, clusterPartition, tsPartition, msgTs, TbMsg.toBytes(msg), msqQueueTtl); |
||||
|
ResultSetFuture resultSetFuture = executeAsyncWrite(boundStatement); |
||||
|
return Futures.transform(resultSetFuture, (Function<ResultSet, Void>) input -> null); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public List<TbMsg> findMsgs(UUID nodeId, long clusterPartition, long tsPartition) { |
||||
|
String select = "SELECT node_id, cluster_partition, ts_partition, ts, msg FROM msg_queue WHERE " + |
||||
|
"node_id = ? AND cluster_partition = ? AND ts_partition = ?"; |
||||
|
PreparedStatement statement = prepare(select); |
||||
|
BoundStatement boundStatement = statement.bind(nodeId, clusterPartition, tsPartition); |
||||
|
ResultSet rows = executeRead(boundStatement); |
||||
|
List<TbMsg> msgs = new ArrayList<>(); |
||||
|
for (Row row : rows) { |
||||
|
msgs.add(TbMsg.fromBytes(row.getBytes("msg"))); |
||||
|
} |
||||
|
return msgs; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -1,2 +1,27 @@ |
|||||
TRUNCATE thingsboard.plugin; |
TRUNCATE thingsboard.plugin; |
||||
TRUNCATE thingsboard.rule; |
TRUNCATE thingsboard.rule; |
||||
|
|
||||
|
-- msg_queue dataset |
||||
|
|
||||
|
INSERT INTO thingsboard.msg_queue (node_id, cluster_partition, ts_partition, ts, msg) |
||||
|
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 200, 201, null); |
||||
|
INSERT INTO thingsboard.msg_queue (node_id, cluster_partition, ts_partition, ts, msg) |
||||
|
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 200, 202, null); |
||||
|
INSERT INTO thingsboard.msg_queue (node_id, cluster_partition, ts_partition, ts, msg) |
||||
|
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 300, 301, null); |
||||
|
|
||||
|
-- ack_queue dataset |
||||
|
INSERT INTO thingsboard.msg_ack_queue (node_id, cluster_partition, ts_partition, msg_id) |
||||
|
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 300, bebaeb60-1888-11e8-bf21-65b5d5335ba9); |
||||
|
INSERT INTO thingsboard.msg_ack_queue (node_id, cluster_partition, ts_partition, msg_id) |
||||
|
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 300, 12baeb60-1888-11e8-bf21-65b5d5335ba9); |
||||
|
INSERT INTO thingsboard.msg_ack_queue (node_id, cluster_partition, ts_partition, msg_id) |
||||
|
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 200, 32baeb60-1888-11e8-bf21-65b5d5335ba9); |
||||
|
|
||||
|
-- processed partition dataset |
||||
|
INSERT INTO thingsboard.processed_msg_partitions (node_id, cluster_partition, ts_partition) |
||||
|
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 100); |
||||
|
INSERT INTO thingsboard.processed_msg_partitions (node_id, cluster_partition, ts_partition) |
||||
|
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 777); |
||||
|
INSERT INTO thingsboard.processed_msg_partitions (node_id, cluster_partition, ts_partition) |
||||
|
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 202, 200); |
||||
@ -1 +1,6 @@ |
|||||
database.type=cassandra |
database.type=cassandra |
||||
|
|
||||
|
cassandra.queue.partitioning=HOURS |
||||
|
cassandra.queue.ack.ttl=1 |
||||
|
cassandra.queue.msg.ttl=1 |
||||
|
cassandra.queue.partitions.ttl=1 |
||||
@ -1,37 +0,0 @@ |
|||||
/** |
|
||||
* 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 lombok.Data; |
|
||||
import org.thingsboard.server.common.data.id.EntityId; |
|
||||
|
|
||||
import java.io.Serializable; |
|
||||
import java.util.UUID; |
|
||||
|
|
||||
/** |
|
||||
* Created by ashvayka on 13.01.18. |
|
||||
*/ |
|
||||
@Data |
|
||||
public final class TbMsg implements Serializable { |
|
||||
|
|
||||
private final UUID id; |
|
||||
private final String type; |
|
||||
private final EntityId originator; |
|
||||
private final TbMsgMetaData metaData; |
|
||||
|
|
||||
private final byte[] data; |
|
||||
|
|
||||
} |
|
||||
@ -1,109 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2017 The Thingsboard Authors |
|
||||
* <p> |
|
||||
* 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 |
|
||||
* <p> |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* <p> |
|
||||
* 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.impl; |
|
||||
|
|
||||
import com.datastax.driver.core.*; |
|
||||
import com.google.common.base.Function; |
|
||||
import com.google.common.util.concurrent.Futures; |
|
||||
import com.google.common.util.concurrent.ListenableFuture; |
|
||||
import com.google.protobuf.ByteString; |
|
||||
import com.google.protobuf.InvalidProtocolBufferException; |
|
||||
import org.springframework.stereotype.Component; |
|
||||
import org.thingsboard.rule.engine.api.TbMsg; |
|
||||
import org.thingsboard.rule.engine.api.TbMsgMetaData; |
|
||||
import org.thingsboard.rule.engine.queue.cassandra.repository.MsgRepository; |
|
||||
import org.thingsboard.rule.engine.queue.cassandra.repository.gen.MsgQueueProtos; |
|
||||
import org.thingsboard.server.common.data.id.EntityId; |
|
||||
import org.thingsboard.server.common.data.id.EntityIdFactory; |
|
||||
|
|
||||
import java.nio.ByteBuffer; |
|
||||
import java.util.ArrayList; |
|
||||
import java.util.List; |
|
||||
import java.util.UUID; |
|
||||
|
|
||||
@Component |
|
||||
public class CassandraMsgRepository extends SimpleAbstractCassandraDao implements MsgRepository { |
|
||||
|
|
||||
private final int msqQueueTtl; |
|
||||
|
|
||||
|
|
||||
public CassandraMsgRepository(Session session, int msqQueueTtl) { |
|
||||
super(session); |
|
||||
this.msqQueueTtl = msqQueueTtl; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<Void> save(TbMsg msg, UUID nodeId, long clusteredHash, long partition, long msgTs) { |
|
||||
String insert = "INSERT INTO msg_queue (node_id, clustered_hash, partition, ts, msg) VALUES (?, ?, ?, ?, ?) USING TTL ?"; |
|
||||
PreparedStatement statement = prepare(insert); |
|
||||
BoundStatement boundStatement = statement.bind(nodeId, clusteredHash, partition, msgTs, toBytes(msg), msqQueueTtl); |
|
||||
ResultSetFuture resultSetFuture = executeAsyncWrite(boundStatement); |
|
||||
return Futures.transform(resultSetFuture, (Function<ResultSet, Void>) input -> null); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public List<TbMsg> findMsgs(UUID nodeId, long clusteredHash, long partition) { |
|
||||
String select = "SELECT node_id, clustered_hash, partition, ts, msg FROM msg_queue WHERE " + |
|
||||
"node_id = ? AND clustered_hash = ? AND partition = ?"; |
|
||||
PreparedStatement statement = prepare(select); |
|
||||
BoundStatement boundStatement = statement.bind(nodeId, clusteredHash, partition); |
|
||||
ResultSet rows = executeRead(boundStatement); |
|
||||
List<TbMsg> msgs = new ArrayList<>(); |
|
||||
for (Row row : rows) { |
|
||||
msgs.add(fromBytes(row.getBytes("msg"))); |
|
||||
} |
|
||||
return msgs; |
|
||||
} |
|
||||
|
|
||||
private ByteBuffer toBytes(TbMsg msg) { |
|
||||
MsgQueueProtos.TbMsgProto.Builder builder = MsgQueueProtos.TbMsgProto.newBuilder(); |
|
||||
builder.setId(msg.getId().toString()); |
|
||||
builder.setType(msg.getType()); |
|
||||
if (msg.getOriginator() != null) { |
|
||||
builder.setEntityType(msg.getOriginator().getEntityType().name()); |
|
||||
builder.setEntityId(msg.getOriginator().getId().toString()); |
|
||||
} |
|
||||
|
|
||||
if (msg.getMetaData() != null) { |
|
||||
MsgQueueProtos.TbMsgProto.TbMsgMetaDataProto.Builder metadataBuilder = MsgQueueProtos.TbMsgProto.TbMsgMetaDataProto.newBuilder(); |
|
||||
metadataBuilder.putAllData(msg.getMetaData().getData()); |
|
||||
builder.addMetaData(metadataBuilder.build()); |
|
||||
} |
|
||||
|
|
||||
builder.setData(ByteString.copyFrom(msg.getData())); |
|
||||
byte[] bytes = builder.build().toByteArray(); |
|
||||
return ByteBuffer.wrap(bytes); |
|
||||
} |
|
||||
|
|
||||
private TbMsg fromBytes(ByteBuffer buffer) { |
|
||||
try { |
|
||||
MsgQueueProtos.TbMsgProto proto = MsgQueueProtos.TbMsgProto.parseFrom(buffer.array()); |
|
||||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|
||||
if (proto.getMetaDataCount() > 0) { |
|
||||
metaData.setData(proto.getMetaData(0).getDataMap()); |
|
||||
} |
|
||||
|
|
||||
EntityId entityId = null; |
|
||||
if (proto.getEntityId() != null) { |
|
||||
entityId = EntityIdFactory.getByTypeAndId(proto.getEntityType(), proto.getEntityId()); |
|
||||
} |
|
||||
|
|
||||
return new TbMsg(UUID.fromString(proto.getId()), proto.getType(), entityId, metaData, proto.getData().toByteArray()); |
|
||||
} catch (InvalidProtocolBufferException e) { |
|
||||
throw new IllegalStateException("Could not parse protobuf for TbMsg", e); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -1,77 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2017 The Thingsboard Authors |
|
||||
* <p> |
|
||||
* 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 |
|
||||
* <p> |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* <p> |
|
||||
* 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.impl; |
|
||||
|
|
||||
import com.datastax.driver.core.*; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.springframework.stereotype.Component; |
|
||||
|
|
||||
import java.util.Map; |
|
||||
import java.util.concurrent.ConcurrentHashMap; |
|
||||
|
|
||||
@Component |
|
||||
@Slf4j |
|
||||
public abstract class SimpleAbstractCassandraDao { |
|
||||
|
|
||||
private ConsistencyLevel defaultReadLevel = ConsistencyLevel.QUORUM; |
|
||||
private ConsistencyLevel defaultWriteLevel = ConsistencyLevel.QUORUM; |
|
||||
private Session session; |
|
||||
private Map<String, PreparedStatement> preparedStatementMap = new ConcurrentHashMap<>(); |
|
||||
|
|
||||
public SimpleAbstractCassandraDao(Session session) { |
|
||||
this.session = session; |
|
||||
} |
|
||||
|
|
||||
protected Session getSession() { |
|
||||
return session; |
|
||||
} |
|
||||
|
|
||||
protected ResultSet executeRead(Statement statement) { |
|
||||
return execute(statement, defaultReadLevel); |
|
||||
} |
|
||||
|
|
||||
protected ResultSet executeWrite(Statement statement) { |
|
||||
return execute(statement, defaultWriteLevel); |
|
||||
} |
|
||||
|
|
||||
protected ResultSetFuture executeAsyncRead(Statement statement) { |
|
||||
return executeAsync(statement, defaultReadLevel); |
|
||||
} |
|
||||
|
|
||||
protected ResultSetFuture executeAsyncWrite(Statement statement) { |
|
||||
return executeAsync(statement, defaultWriteLevel); |
|
||||
} |
|
||||
|
|
||||
protected PreparedStatement prepare(String query) { |
|
||||
return preparedStatementMap.computeIfAbsent(query, i -> getSession().prepare(i)); |
|
||||
} |
|
||||
|
|
||||
private ResultSet execute(Statement statement, ConsistencyLevel level) { |
|
||||
log.debug("Execute cassandra statement {}", statement); |
|
||||
if (statement.getConsistencyLevel() == null) { |
|
||||
statement.setConsistencyLevel(level); |
|
||||
} |
|
||||
return getSession().execute(statement); |
|
||||
} |
|
||||
|
|
||||
private ResultSetFuture executeAsync(Statement statement, ConsistencyLevel level) { |
|
||||
log.debug("Execute cassandra async statement {}", statement); |
|
||||
if (statement.getConsistencyLevel() == null) { |
|
||||
statement.setConsistencyLevel(level); |
|
||||
} |
|
||||
return getSession().executeAsync(statement); |
|
||||
} |
|
||||
} |
|
||||
@ -1,48 +0,0 @@ |
|||||
/** |
|
||||
* 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.junit.Before; |
|
||||
import org.junit.Test; |
|
||||
import org.mockito.Mock; |
|
||||
import org.thingsboard.rule.engine.queue.cassandra.repository.AckRepository; |
|
||||
import org.thingsboard.rule.engine.queue.cassandra.repository.MsgRepository; |
|
||||
|
|
||||
public class CassandraMsqQueueTest { |
|
||||
|
|
||||
private CassandraMsqQueue msqQueue; |
|
||||
|
|
||||
@Mock |
|
||||
private MsgRepository msgRepository; |
|
||||
@Mock |
|
||||
private AckRepository ackRepository; |
|
||||
@Mock |
|
||||
private UnprocessedMsgFilter unprocessedMsgFilter; |
|
||||
@Mock |
|
||||
private QueuePartitioner queuePartitioner; |
|
||||
|
|
||||
@Before |
|
||||
public void init() { |
|
||||
msqQueue = new CassandraMsqQueue(msgRepository, ackRepository, unprocessedMsgFilter, queuePartitioner); |
|
||||
} |
|
||||
|
|
||||
@Test |
|
||||
public void msgCanBeSaved() { |
|
||||
// todo-vp: implement
|
|
||||
} |
|
||||
|
|
||||
|
|
||||
} |
|
||||
@ -1,30 +0,0 @@ |
|||||
/** |
|
||||
* 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.impl; |
|
||||
|
|
||||
import org.cassandraunit.CassandraCQLUnit; |
|
||||
import org.cassandraunit.dataset.cql.ClassPathCQLDataSet; |
|
||||
import org.junit.ClassRule; |
|
||||
|
|
||||
|
|
||||
public abstract class SimpleAbstractCassandraDaoTest { |
|
||||
|
|
||||
@ClassRule |
|
||||
public static CassandraCQLUnit cassandraUnit = new CassandraCQLUnit( |
|
||||
new ClassPathCQLDataSet("cassandra/system-test.cql", "thingsboard")); |
|
||||
|
|
||||
|
|
||||
} |
|
||||
@ -1,75 +0,0 @@ |
|||||
CREATE TABLE IF NOT EXISTS thingsboard.msg_queue ( |
|
||||
node_id timeuuid, |
|
||||
clustered_hash bigint, |
|
||||
partition bigint, |
|
||||
ts bigint, |
|
||||
msg blob, |
|
||||
PRIMARY KEY ((node_id, clustered_hash, partition), ts)) |
|
||||
WITH CLUSTERING ORDER BY (ts DESC) |
|
||||
AND compaction = { |
|
||||
'class': 'org.apache.cassandra.db.compaction.DateTieredCompactionStrategy', |
|
||||
'min_threshold': '5', |
|
||||
'base_time_seconds': '43200', |
|
||||
'max_window_size_seconds': '43200', |
|
||||
'tombstone_threshold': '0.9', |
|
||||
'unchecked_tombstone_compaction': 'true' |
|
||||
}; |
|
||||
|
|
||||
|
|
||||
CREATE TABLE IF NOT EXISTS thingsboard.msg_ack_queue ( |
|
||||
node_id timeuuid, |
|
||||
clustered_hash bigint, |
|
||||
partition bigint, |
|
||||
msg_id timeuuid, |
|
||||
PRIMARY KEY ((node_id, clustered_hash, partition), msg_id)) |
|
||||
WITH CLUSTERING ORDER BY (msg_id DESC) |
|
||||
AND compaction = { |
|
||||
'class': 'org.apache.cassandra.db.compaction.DateTieredCompactionStrategy', |
|
||||
'min_threshold': '5', |
|
||||
'base_time_seconds': '43200', |
|
||||
'max_window_size_seconds': '43200', |
|
||||
'tombstone_threshold': '0.9', |
|
||||
'unchecked_tombstone_compaction': 'true' |
|
||||
}; |
|
||||
|
|
||||
CREATE TABLE IF NOT EXISTS thingsboard.processed_msg_partitions ( |
|
||||
node_id timeuuid, |
|
||||
clustered_hash bigint, |
|
||||
partition bigint, |
|
||||
PRIMARY KEY ((node_id, clustered_hash), partition)) |
|
||||
WITH CLUSTERING ORDER BY (partition DESC) |
|
||||
AND compaction = { |
|
||||
'class': 'org.apache.cassandra.db.compaction.DateTieredCompactionStrategy', |
|
||||
'min_threshold': '5', |
|
||||
'base_time_seconds': '43200', |
|
||||
'max_window_size_seconds': '43200', |
|
||||
'tombstone_threshold': '0.9', |
|
||||
'unchecked_tombstone_compaction': 'true' |
|
||||
}; |
|
||||
|
|
||||
|
|
||||
|
|
||||
-- msg_queue dataset |
|
||||
|
|
||||
INSERT INTO thingsboard.msg_queue (node_id, clustered_hash, partition, ts, msg) |
|
||||
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 200, 201, null); |
|
||||
INSERT INTO thingsboard.msg_queue (node_id, clustered_hash, partition, ts, msg) |
|
||||
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 200, 202, null); |
|
||||
INSERT INTO thingsboard.msg_queue (node_id, clustered_hash, partition, ts, msg) |
|
||||
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 300, 301, null); |
|
||||
|
|
||||
-- ack_queue dataset |
|
||||
INSERT INTO msg_ack_queue (node_id, clustered_hash, partition, msg_id) |
|
||||
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 300, bebaeb60-1888-11e8-bf21-65b5d5335ba9); |
|
||||
INSERT INTO msg_ack_queue (node_id, clustered_hash, partition, msg_id) |
|
||||
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 300, 12baeb60-1888-11e8-bf21-65b5d5335ba9); |
|
||||
INSERT INTO msg_ack_queue (node_id, clustered_hash, partition, msg_id) |
|
||||
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 200, 32baeb60-1888-11e8-bf21-65b5d5335ba9); |
|
||||
|
|
||||
-- processed partition dataset |
|
||||
INSERT INTO processed_msg_partitions (node_id, clustered_hash, partition) |
|
||||
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 100); |
|
||||
INSERT INTO processed_msg_partitions (node_id, clustered_hash, partition) |
|
||||
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 101, 777); |
|
||||
INSERT INTO processed_msg_partitions (node_id, clustered_hash, partition) |
|
||||
VALUES (055eee50-1883-11e8-b380-65b5d5335ba9, 202, 200); |
|
||||
Loading…
Reference in new issue