2  * ================================================================================
 
   3  * Copyright (c) 2018 AT&T Intellectual Property. All rights reserved.
 
   4  * ================================================================================
 
   5  * Licensed under the Apache License, Version 2.0 (the "License");
 
   6  * you may not use this file except in compliance with the License.
 
   7  * You may obtain a copy of the License at
 
   9  *      http://www.apache.org/licenses/LICENSE-2.0
 
  11  * Unless required by applicable law or agreed to in writing, software
 
  12  * distributed under the License is distributed on an "AS IS" BASIS,
 
  13  * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
 
  14  * See the License for the specific language governing permissions and
 
  15  * limitations under the License.
 
  16  * ============LICENSE_END=========================================================
 
  20 package org.onap.dcae.analytics.web.dmaap;
 
  22 import static org.apache.commons.text.StringEscapeUtils.unescapeJava;
 
  23 import static org.apache.commons.text.StringEscapeUtils.unescapeJson;
 
  24 import static org.onap.dcae.analytics.model.AnalyticsHttpConstants.REQUEST_ID_HEADER_KEY;
 
  25 import static org.onap.dcae.analytics.model.AnalyticsModelConstants.ANALYTICS_REQUEST_ID_DELIMITER;
 
  27 import java.io.IOException;
 
  28 import java.util.Collections;
 
  29 import java.util.Date;
 
  30 import java.util.LinkedList;
 
  31 import java.util.List;
 
  32 import java.util.stream.IntStream;
 
  34 import javax.annotation.Nonnull;
 
  35 import javax.annotation.Nullable;
 
  37 import org.onap.dcae.analytics.model.AnalyticsHttpConstants;
 
  38 import org.onap.dcae.analytics.model.DmaapMrConstants;
 
  39 import org.onap.dcae.analytics.tca.core.exception.AnalyticsParsingException;
 
  40 import org.onap.dcae.analytics.tca.core.util.LogSpec;
 
  41 import org.onap.dcae.analytics.web.util.AnalyticsHttpUtils;
 
  42 import org.onap.dcae.utils.eelf.logger.api.log.EELFLogFactory;
 
  43 import org.onap.dcae.utils.eelf.logger.api.log.EELFLogger;
 
  44 import org.onap.dcae.utils.eelf.logger.api.spec.AuditLogSpec;
 
  45 import org.onap.dcae.utils.eelf.logger.api.spec.ErrorLogSpec;
 
  46 import org.springframework.integration.splitter.AbstractMessageSplitter;
 
  47 import org.springframework.integration.support.MessageBuilder;
 
  48 import org.springframework.messaging.Message;
 
  50 import com.fasterxml.jackson.databind.JsonNode;
 
  51 import com.fasterxml.jackson.databind.ObjectMapper;
 
  54  * DMaaP MR message splitter split the incoming messages into batch of given batch size
 
  56  * @author Rajiv Singla
 
  58 public class MrMessageSplitter extends AbstractMessageSplitter {
 
  60     private static final EELFLogger eelfLogger = EELFLogFactory.getLogger(MrMessageSplitter.class);
 
  62     private final ObjectMapper objectMapper;
 
  63     private final Integer batchSize;
 
  65     public MrMessageSplitter(@Nonnull final ObjectMapper objectMapper,
 
  66                              @Nonnull final Integer batchSize) {
 
  67         this.objectMapper = objectMapper;
 
  68         this.batchSize = batchSize;
 
  72     protected Object splitMessage(final Message<?> message) {
 
  74         final String requestId = AnalyticsHttpUtils.getRequestId(message.getHeaders());
 
  75         final List<String> dmaapMessages = convertJsonToStringMessages(requestId, String.class.cast(message.getPayload()).trim());
 
  77         final String transactionId = AnalyticsHttpUtils.getTransactionId(message.getHeaders());
 
  79         final Date requestBeginTimestamp = AnalyticsHttpUtils.getTimestampFromHeaders(message.getHeaders(),
 
  80                 AnalyticsHttpConstants.REQUEST_BEGIN_TS_HEADER_KEY);
 
  81         final AuditLogSpec auditLogSpec = LogSpec.createAuditLogSpec(requestId, requestBeginTimestamp);
 
  83         eelfLogger.auditLog().info("Request Id: {}, Transaction Id: {}, dmaapMessages: {},"
 
  84                 + " Received new messages from DMaaP MR. Count: {}",
 
  85                 auditLogSpec, requestId, transactionId, dmaapMessages.toString(), String.valueOf(dmaapMessages.size()));
 
  87         final List<List<String>> messagePartitions = partition(dmaapMessages, batchSize);
 
  88         eelfLogger.auditLog().info("Request Id: {}, Transaction Id: {},  Max allowed messages per batch: {}. " +
 
  89                 "No of batches created: {}",
 
  90                 auditLogSpec, requestId, transactionId, String.valueOf(batchSize), String.valueOf(messagePartitions.size()));
 
  92         // append batch id to request id header
 
  93         return messagePartitions.isEmpty() ? null : IntStream.range(0, messagePartitions.size())
 
  94                 .mapToObj(batchIndex ->
 
  96                                 .withPayload(messagePartitions.get(0))
 
  97                                 .copyHeaders(message.getHeaders())
 
  98                                 .setHeader(REQUEST_ID_HEADER_KEY,
 
  99                                         requestId + ANALYTICS_REQUEST_ID_DELIMITER + batchIndex)
 
 106      * Converts DMaaP MR subscriber messages json string to List of messages. If message Json String is empty
 
 109      * @param messagesJsonString json messages String
 
 111      * @return List containing DMaaP MR Messages
 
 113     private List<String> convertJsonToStringMessages(String requestId, @Nullable final String messagesJsonString) {
 
 115         final LinkedList<String> messages = new LinkedList<>();
 
 117         // If message string is not null or not empty parse json message array to List of string messages
 
 118         if (messagesJsonString != null && !messagesJsonString.trim().isEmpty()
 
 119                 && !DmaapMrConstants.SUBSCRIBER_EMPTY_MESSAGE_RESPONSE_STRING.equals(messagesJsonString.trim())) {
 
 123                 final JsonNode rootNode = objectMapper.readTree(messagesJsonString);
 
 124                 // iterate over root node and parse arrays messages
 
 125                 for (JsonNode jsonNode : rootNode) {
 
 126                     // if array parse it is array of messages
 
 127                     final String incomingMessageString = jsonNode.toString();
 
 128                     if (jsonNode.isArray()) {
 
 129                         final List messageList = objectMapper.readValue(incomingMessageString, List.class);
 
 130                         for (Object message : messageList) {
 
 131                             final String jsonMessageString = objectMapper.writeValueAsString(message);
 
 132                             addUnescapedJsonToMessage(messages, jsonMessageString);
 
 135                         // parse it as object
 
 136                         addUnescapedJsonToMessage(messages, incomingMessageString);
 
 140             } catch (IOException e) {
 
 141                 ErrorLogSpec errorLogSpec = LogSpec.createErrorLogSpec(requestId);
 
 142                 eelfLogger.errorLog().error("Unable to convert subscriber Json String to Messages. " +
 
 143                         "Subscriber Response String: {}, Json Error: {}", errorLogSpec, messagesJsonString, e.toString());
 
 144                 String errorMessage = String.format("Unable to convert subscriber Json String to Messages. " +
 
 145                         "Subscriber Response String: %s, Json Error: %s", messagesJsonString, e);
 
 146                 throw new AnalyticsParsingException(errorMessage, e);
 
 154      * Adds unescaped Json messages to given messages list
 
 156      * @param messages message list in which unescaped messages will be added
 
 157      * @param incomingMessageString incoming message string that may need to be escaped
 
 159     private static void addUnescapedJsonToMessage(List<String> messages, String incomingMessageString) {
 
 160         if (incomingMessageString.startsWith("\"") && incomingMessageString.endsWith("\"")) {
 
 161             messages.add(unescapeJava(unescapeJson(
 
 162                     incomingMessageString.substring(1, incomingMessageString.length() - 1))));
 
 164             messages.add(unescapeJava(unescapeJson(incomingMessageString)));
 
 169      * Partition list into multiple lists
 
 171      * @param list input list that needs to be broken into chunks
 
 172      * @param batchSize batch size for each list
 
 173      * @param <E> element type of the list
 
 175      * @return List containing list of entries of specified batch size
 
 177     private static <E> List<List<E>> partition(List<E> list, final Integer batchSize) {
 
 179         if (list == null || batchSize == null || batchSize <= 0 || list.size() < batchSize) {
 
 180             return Collections.singletonList(list);
 
 183         final List<List<E>> result = new LinkedList<>();
 
 185         for (int i = 0; i < list.size(); i++) {
 
 187             if (i == 0 || i % batchSize == 0) {
 
 188                 List<E> sublist = new LinkedList<>();
 
 192             final List<E> lastSubList = result.get(result.size() - 1);
 
 193             lastSubList.add(list.get(i));