DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
dunedaq::kafkaopmon::OpMonPublisher Class Reference

#include <OpMonPublisher.hpp>

Public Member Functions

 OpMonPublisher (const nlohmann::json &conf)
 OpMonPublisher ()=delete
 OpMonPublisher (const OpMonPublisher &)=delete
OpMonPublisher & operator= (const OpMonPublisher &)=delete
 OpMonPublisher (OpMonPublisher &&)=delete
OpMonPublisher & operator= (OpMonPublisher &&)=delete
 ~OpMonPublisher ()
void publish (dunedaq::opmon::OpMonEntry &&) const

Protected Member Functions

std::string extract_topic (const dunedaq::opmon::OpMonEntry &) const noexcept
std::string extract_key (const dunedaq::opmon::OpMonEntry &e) const noexcept

Private Attributes

std::unique_ptr< RdKafka::Producer > m_producer
std::string m_default_topic = "monitoring.opmon_stream"

Detailed Description

Definition at line 55 of file OpMonPublisher.hpp.

Constructor & Destructor Documentation

◆ OpMonPublisher() [1/4]

OpMonPublisher::OpMonPublisher ( const nlohmann::json & conf)

Definition at line 13 of file OpMonPublisher.cpp.

14{
15
16 RdKafka::Conf* k_conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
17 std::string errstr;
18
19 // Good observations on threadsafety here
20 // https://www.confluent.io/blog/modern-cpp-kafka-api-for-safe-easy-messaging/
21
22 auto it = conf.find("bootstrap");
23 if (it == conf.end()) {
24 throw MissingParameter(ERS_HERE, "bootstrap", nlohmann::to_string(conf));
25 }
26
27 k_conf->set("bootstrap.servers", *it, errstr);
28 if (!errstr.empty()) {
29 throw FailedConfiguration(ERS_HERE, "bootstrap.servers", errstr);
30 }
31
32 std::string client_id;
33 it = conf.find("cliend_id");
34 if (it != conf.end())
35 client_id = *it;
36 else
37 client_id = "kafkaopmon_default_producer";
38
39 k_conf->set("client.id", client_id, errstr);
40 if (!errstr.empty()) {
41 ers::error(FailedConfiguration(ERS_HERE, "client.id", errstr));
42 }
43
44 // Create producer instance
45 m_producer.reset(RdKafka::Producer::create(k_conf, errstr));
46
47 if (!m_producer) {
48 throw FailedProducerCreation(ERS_HERE, errstr);
49 }
50
51 it = conf.find("default_topic");
52 if (it != conf.end())
53 m_default_topic = "monitoring." + it->get<std::string>();
54}
#define ERS_HERE
std::unique_ptr< RdKafka::Producer > m_producer
void error(const Issue &issue)
Definition ers.hpp:101

◆ OpMonPublisher() [2/4]

dunedaq::kafkaopmon::OpMonPublisher::OpMonPublisher ( )
delete

◆ OpMonPublisher() [3/4]

dunedaq::kafkaopmon::OpMonPublisher::OpMonPublisher ( const OpMonPublisher & )
delete

◆ OpMonPublisher() [4/4]

dunedaq::kafkaopmon::OpMonPublisher::OpMonPublisher ( OpMonPublisher && )
delete

◆ ~OpMonPublisher()

OpMonPublisher::~OpMonPublisher ( )

Definition at line 56 of file OpMonPublisher.cpp.

57{
58
59 int timeout_ms = 500;
60 RdKafka::ErrorCode err = m_producer->flush(timeout_ms);
61
62 if (err == RdKafka::ERR__TIMED_OUT) {
63 ers::warning(TimeoutReachedWhileFlushing(ERS_HERE, timeout_ms));
64 }
65}
void warning(const Issue &issue)
Definition ers.hpp:150

Member Function Documentation

◆ extract_key()

std::string dunedaq::kafkaopmon::OpMonPublisher::extract_key ( const dunedaq::opmon::OpMonEntry & e) const
inlineprotectednoexcept

Definition at line 73 of file OpMonPublisher.hpp.

74 {
75 return dunedaq::opmonlib::to_string(e.origin()) + '/' + e.measurement();
76 }
const ::dunedaq::opmon::OpMonId & origin() const
const std::string & measurement() const
std::string to_string(const dunedaq::opmon::OpMonId &)
Definition Utils.cpp:166

◆ extract_topic()

std::string dunedaq::kafkaopmon::OpMonPublisher::extract_topic ( const dunedaq::opmon::OpMonEntry & ) const
inlineprotectednoexcept

Definition at line 72 of file OpMonPublisher.hpp.

72{ return m_default_topic; }

◆ operator=() [1/2]

OpMonPublisher & dunedaq::kafkaopmon::OpMonPublisher::operator= ( const OpMonPublisher & )
delete

◆ operator=() [2/2]

OpMonPublisher & dunedaq::kafkaopmon::OpMonPublisher::operator= ( OpMonPublisher && )
delete

◆ publish()

void OpMonPublisher::publish ( dunedaq::opmon::OpMonEntry && entry) const

Definition at line 68 of file OpMonPublisher.cpp.

69{
70
71 std::string binary;
72 entry.SerializeToString(&binary);
73
74 auto topic = extract_topic(entry);
75 auto key = extract_key(entry);
76
77 RdKafka::ErrorCode err = m_producer->produce(topic,
78 RdKafka::Topic::PARTITION_UA,
79 RdKafka::Producer::RK_MSG_COPY,
80 const_cast<char*>(binary.c_str()),
81 binary.size(),
82 key.c_str(),
83 key.size(),
84 0,
85 nullptr);
86
87 if (err == RdKafka::ERR_NO_ERROR)
88 return;
89
90 std::string err_cause;
91
92 switch (err) {
93 case RdKafka::ERR__QUEUE_FULL:
94 err_cause = "maximum number of outstanding messages reached";
95 break;
96 case RdKafka::ERR_MSG_SIZE_TOO_LARGE:
97 err_cause = "message too large";
98 break;
99 case RdKafka::ERR__UNKNOWN_PARTITION:
100 err_cause = "Unknown partition";
101 break;
102 case RdKafka::ERR__UNKNOWN_TOPIC:
103 err_cause = "Unknown topic (";
104 err_cause += topic;
105 err_cause += ')';
106 break;
107 default:
108 err_cause = "unknown";
109 break;
110 }
111
112 throw FailedProduce(ERS_HERE, key, err_cause);
113}
std::string extract_key(const dunedaq::opmon::OpMonEntry &e) const noexcept
std::string extract_topic(const dunedaq::opmon::OpMonEntry &) const noexcept

Member Data Documentation

◆ m_default_topic

std::string dunedaq::kafkaopmon::OpMonPublisher::m_default_topic = "monitoring.opmon_stream"
private

Definition at line 80 of file OpMonPublisher.hpp.

◆ m_producer

std::unique_ptr<RdKafka::Producer> dunedaq::kafkaopmon::OpMonPublisher::m_producer
private

Definition at line 79 of file OpMonPublisher.hpp.


The documentation for this class was generated from the following files: