2 * Copyright © 2019 IBM.
3 * Modifications Copyright © 2018-2019 AT&T Intellectual Property.
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.
18 package org.onap.ccsdk.cds.blueprintsprocessor.message.service
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.ProducerConfig.ACKS_CONFIG
24 import org.apache.kafka.clients.producer.ProducerConfig.BOOTSTRAP_SERVERS_CONFIG
25 import org.apache.kafka.clients.producer.ProducerConfig.CLIENT_ID_CONFIG
26 import org.apache.kafka.clients.producer.ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG
27 import org.apache.kafka.clients.producer.ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG
28 import org.apache.kafka.clients.producer.ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG
29 import org.apache.kafka.clients.producer.ProducerRecord
30 import org.apache.kafka.common.header.internals.RecordHeader
31 import org.apache.kafka.common.serialization.ByteArraySerializer
32 import org.apache.kafka.common.serialization.StringSerializer
33 import org.onap.ccsdk.cds.blueprintsprocessor.message.KafkaBasicAuthMessageProducerProperties
34 import org.onap.ccsdk.cds.controllerblueprints.core.asJsonString
35 import org.onap.ccsdk.cds.controllerblueprints.core.defaultToUUID
36 import org.slf4j.LoggerFactory
37 import java.nio.charset.Charset
39 class KafkaBasicAuthMessageProducerService(
40 private val messageProducerProperties: KafkaBasicAuthMessageProducerProperties
42 BlueprintMessageProducerService {
44 private val log = LoggerFactory.getLogger(KafkaBasicAuthMessageProducerService::class.java)!!
46 private var kafkaProducer: KafkaProducer<String, ByteArray>? = null
48 private val messageLoggerService = MessageLoggerService()
50 override suspend fun sendMessageNB(message: Any): Boolean {
51 checkNotNull(messageProducerProperties.topic) { "default topic is not configured" }
52 return sendMessageNB(messageProducerProperties.topic!!, message)
55 override suspend fun sendMessageNB(message: Any, headers: MutableMap<String, String>?): Boolean {
56 checkNotNull(messageProducerProperties.topic) { "default topic is not configured" }
57 return sendMessageNB(messageProducerProperties.topic!!, message, headers)
60 override suspend fun sendMessageNB(
63 headers: MutableMap<String, String>?
65 val byteArrayMessage = when (message) {
66 is String -> message.toByteArray(Charset.defaultCharset())
67 else -> message.asJsonString().toByteArray(Charset.defaultCharset())
70 val record = ProducerRecord<String, ByteArray>(topic, defaultToUUID(), byteArrayMessage)
71 val recordHeaders = record.headers()
72 messageLoggerService.messageProducing(recordHeaders)
74 headers.forEach { (key, value) -> recordHeaders.add(RecordHeader(key, value.toByteArray())) }
76 val callback = Callback { metadata, exception ->
77 log.trace("message published to(${metadata.topic()}), offset(${metadata.offset()}), headers :$headers")
79 messageTemplate().send(record, callback)
83 fun messageTemplate(additionalConfig: Map<String, ByteArray>? = null): KafkaProducer<String, ByteArray> {
84 log.trace("Client Properties : ${ToStringBuilder.reflectionToString(messageProducerProperties)}")
85 val configProps = hashMapOf<String, Any>()
86 configProps[BOOTSTRAP_SERVERS_CONFIG] = messageProducerProperties.bootstrapServers
87 configProps[KEY_SERIALIZER_CLASS_CONFIG] = StringSerializer::class.java
88 configProps[VALUE_SERIALIZER_CLASS_CONFIG] = ByteArraySerializer::class.java
89 configProps[ACKS_CONFIG] = messageProducerProperties.acks
90 configProps[ENABLE_IDEMPOTENCE_CONFIG] = messageProducerProperties.enableIdempotence
91 if (messageProducerProperties.clientId != null) {
92 configProps[CLIENT_ID_CONFIG] = messageProducerProperties.clientId!!
94 // TODO("Security Implementation based on type")
96 // Add additional Properties
97 if (additionalConfig != null) {
98 configProps.putAll(additionalConfig)
101 if (kafkaProducer == null) {
102 kafkaProducer = KafkaProducer(configProps)
104 return kafkaProducer!!