DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
OpMonPublisher.cpp
Go to the documentation of this file.
1
8
10
11using namespace dunedaq::kafkaopmon;
12
13OpMonPublisher::OpMonPublisher(const nlohmann::json& conf)
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}
55
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}
66
67void
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}
#define ERS_HERE
std::unique_ptr< RdKafka::Producer > m_producer
std::string extract_key(const dunedaq::opmon::OpMonEntry &e) const noexcept
std::string extract_topic(const dunedaq::opmon::OpMonEntry &) const noexcept
void publish(dunedaq::opmon::OpMonEntry &&) const
void warning(const Issue &issue)
Definition ers.hpp:150
void error(const Issue &issue)
Definition ers.hpp:101