931f052ed54b6d180ef0f1525b59bcdbaaf028b1
[ccsdk/cds.git] /
1 /*
2  *  Copyright © 2019 IBM.
3  *  Modifications Copyright © 2018-2019 AT&T Intellectual Property.
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
8  *
9  *      http://www.apache.org/licenses/LICENSE-2.0
10  *
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  */
17
18 package org.onap.ccsdk.cds.blueprintsprocessor.message.service
19
20 import org.apache.commons.lang.builder.ToStringBuilder
21 import org.apache.kafka.clients.producer.Callback
22 import org.apache.kafka.clients.producer.KafkaProducer
23 import org.apache.kafka.clients.producer.ProducerRecord
24 import org.apache.kafka.common.header.internals.RecordHeader
25 import org.onap.ccsdk.cds.blueprintsprocessor.message.MessageProducerProperties
26 import org.onap.ccsdk.cds.controllerblueprints.core.asJsonString
27 import org.onap.ccsdk.cds.controllerblueprints.core.defaultToUUID
28 import org.slf4j.LoggerFactory
29 import java.nio.charset.Charset
30
31 class KafkaMessageProducerService(
32     private val messageProducerProperties: MessageProducerProperties
33 ) :
34     BlueprintMessageProducerService {
35
36     private val log = LoggerFactory.getLogger(KafkaMessageProducerService::class.java)!!
37
38     private var kafkaProducer: KafkaProducer<String, ByteArray>? = null
39
40     private val messageLoggerService = MessageLoggerService()
41
42     override suspend fun sendMessageNB(message: Any): Boolean {
43         checkNotNull(messageProducerProperties.topic) { "default topic is not configured" }
44         return sendMessageNB(messageProducerProperties.topic!!, message)
45     }
46
47     override suspend fun sendMessageNB(message: Any, headers: MutableMap<String, String>?): Boolean {
48         checkNotNull(messageProducerProperties.topic) { "default topic is not configured" }
49         return sendMessageNB(messageProducerProperties.topic!!, message, headers)
50     }
51
52     override suspend fun sendMessageNB(
53         topic: String,
54         message: Any,
55         headers: MutableMap<String, String>?
56     ): Boolean {
57         val byteArrayMessage = when (message) {
58             is String -> message.toByteArray(Charset.defaultCharset())
59             else -> message.asJsonString().toByteArray(Charset.defaultCharset())
60         }
61
62         val record = ProducerRecord<String, ByteArray>(topic, defaultToUUID(), byteArrayMessage)
63         val recordHeaders = record.headers()
64         messageLoggerService.messageProducing(recordHeaders)
65         headers?.let {
66             headers.forEach { (key, value) -> recordHeaders.add(RecordHeader(key, value.toByteArray())) }
67         }
68         val callback = Callback { metadata, exception ->
69             log.trace("message published to(${metadata.topic()}), offset(${metadata.offset()}), headers :$headers")
70         }
71         messageTemplate().send(record, callback)
72         return true
73     }
74
75     fun messageTemplate(additionalConfig: Map<String, ByteArray>? = null): KafkaProducer<String, ByteArray> {
76         log.trace("Producer client properties : ${ToStringBuilder.reflectionToString(messageProducerProperties)}")
77         val configProps = messageProducerProperties.getConfig()
78
79         /** Add additional Properties */
80         if (additionalConfig != null)
81             configProps.putAll(additionalConfig)
82
83         if (kafkaProducer == null)
84             kafkaProducer = KafkaProducer(configProps)
85
86         return kafkaProducer!!
87     }
88 }