summaryrefslogtreecommitdiffstats
path: root/components/kpi-computation-ms/src/main/java/org/onap/dcaegen2/kpi/dmaap/NotificationProducer.java
blob: c5be6cc023e6ec1aea468a2d7a044b3c04c4aa5e (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
/*-
 * ============LICENSE_START=======================================================
 *  Copyright (C) 2021 China Mobile.
 *  Copyright (C) 2022 Wipro Limited.
 * ================================================================================
 * 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.
 *
 * SPDX-License-Identifier: Apache-2.0
 * ============LICENSE_END=========================================================
 */

package org.onap.dcaegen2.kpi.dmaap;

import java.io.IOException;
import org.onap.dcaegen2.services.sdk.rest.services.dmaap.client.api.MessageRouterPublisher;
import org.onap.dcaegen2.services.sdk.rest.services.dmaap.client.model.MessageRouterPublishRequest;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.google.gson.JsonPrimitive;
import reactor.core.publisher.Flux;

/**
 * Produces Notification on DMAAP events.
 */
public class NotificationProducer {
    private static Logger logger = LoggerFactory.getLogger(NotificationProducer.class);
    private MessageRouterPublisher messageRouterPublisher;
    private MessageRouterPublishRequest messageRouterPublishRequest;
    
    /**
     * Parameterized constructor.
     */
    public NotificationProducer(MessageRouterPublisher messageRouterPublisher, MessageRouterPublishRequest messageRouterPublishRequest) {
        super();
        this.messageRouterPublisher = messageRouterPublisher;
        this.messageRouterPublishRequest = messageRouterPublishRequest;
    }

    /**
     * sends notification to dmaap.
     */
    public void sendNotification(String msg) throws IOException {
        Flux.just(1, 2, 3)
           .map(JsonPrimitive::new)
           .transform(input -> messageRouterPublisher.put(messageRouterPublishRequest, input))
           .subscribe(resp -> {
               if (resp.successful()) {
                   logger.debug("Sent a batch of messages to the MR");
                       } else {
                           logger.warn("Message sending has failed: {}", resp.failReason());
                     }
                },
                ex -> {
                    logger.warn("An unexpected error while sending messages to DMaaP", ex);
                 });
    }
}