Browse Source

refactor code

pull/7103/head
Yuriy Lytvynchuk 4 years ago
parent
commit
6d6e7b212c
  1. 25
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNode.java

25
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbSplitArrayMsgNode.java

@ -29,8 +29,6 @@ import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
@Slf4j @Slf4j
@ -42,6 +40,7 @@ import java.util.concurrent.ExecutionException;
nodeDetails = "Split the array fetched from the msg body. If the msg data is not a JSON object returns the " nodeDetails = "Split the array fetched from the msg body. If the msg data is not a JSON object returns the "
+ "incoming message as outbound message with <code>Failure</code> chain, otherwise returns " + "incoming message as outbound message with <code>Failure</code> chain, otherwise returns "
+ "inner objects of the extracted array as separate messages via <code>Success</code> chain.", + "inner objects of the extracted array as separate messages via <code>Success</code> chain.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
icon = "content_copy", icon = "content_copy",
configDirective = "tbNodeEmptyConfig" configDirective = "tbNodeEmptyConfig"
) )
@ -59,31 +58,19 @@ public class TbSplitArrayMsgNode implements TbNode {
JsonNode jsonNode = JacksonUtil.toJsonNode(msg.getData()); JsonNode jsonNode = JacksonUtil.toJsonNode(msg.getData());
if (jsonNode.isArray()) { if (jsonNode.isArray()) {
ArrayNode data = (ArrayNode) jsonNode; ArrayNode data = (ArrayNode) jsonNode;
List<TbMsg> messages = new ArrayList<>(); if (data.size() == 1) {
data.forEach(msgNode -> { ctx.tellSuccess(TbMsg.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), JacksonUtil.toString(data.get(0))));
messages.add(createMsg(msg, msgNode, data.size() > 1));
});
if (messages.size() == 1) {
ctx.tellSuccess(messages.get(0));
} else { } else {
ctx.ack(msg); ctx.ack(msg);
for (TbMsg newMsg : messages) { data.forEach(msgNode -> {
ctx.tellSuccess(newMsg); ctx.tellSuccess(TbMsg.newMsg(msg.getQueueName(), msg.getType(), msg.getOriginator(), msg.getMetaData(), JacksonUtil.toString(msgNode)));
} });
} }
} else { } else {
ctx.tellFailure(msg, new RuntimeException("Msg data is not a JSON Array!")); ctx.tellFailure(msg, new RuntimeException("Msg data is not a JSON Array!"));
} }
} }
private TbMsg createMsg(TbMsg msg, JsonNode msgNode, boolean newMessage) {
if (newMessage) {
return TbMsg.newMsg(msg.getQueueName(), msg.getType(), msg.getOriginator(), msg.getMetaData(), JacksonUtil.toString(msgNode));
} else {
return TbMsg.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), JacksonUtil.toString(msgNode));
}
}
@Override @Override
public void destroy() { public void destroy() {

Loading…
Cancel
Save