KProducer.scala 840 Bytes
Newer Older
Sebastian Schüpbach's avatar
Sebastian Schüpbach committed
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
package ch.memobase

import ch.memobase.models.DeleteMessage
import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord}
import org.apache.logging.log4j.scala.Logging


abstract class KProducer {
  self: AppSettings with Logging =>

  private lazy val producer = new KafkaProducer[String, String](producerProps)

  def sendDelete(id: DeleteMessage): Unit = {
    logger.info(s"Sending delete command for $id")
    val producerRecord = new ProducerRecord[String, String](id.recordId, null)
    producerRecord.headers()
      .add("recordSetId", id.recordSetId.getBytes())
      .add("institutionId", id.institutionId.getBytes())
      .add("sessionId", id.sessionId.getBytes())
    producer.send(producerRecord)
  }

  def closeProducer(): Unit = {
    logger.debug("Closing Kafka producer instance")
    producer.close()
  }
}