DUNE-DAQ
DUNE Trigger and Data Acquisition software
Toggle main menu visibility
Loading...
Searching...
No Matches
dunedaq
sourcecode
kafkaopmon
src
OpMonPublisher.cpp
Go to the documentation of this file.
1
8
9
#include "
kafkaopmon/OpMonPublisher.hpp
"
10
11
using namespace
dunedaq::kafkaopmon
;
12
13
OpMonPublisher::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
56
OpMonPublisher::~OpMonPublisher
()
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
67
void
68
OpMonPublisher::publish
(
dunedaq::opmon::OpMonEntry
&& entry)
const
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
}
ERS_HERE
#define ERS_HERE
Definition
LocalContext.hpp:127
OpMonPublisher.hpp
dunedaq::kafkaopmon::OpMonPublisher::m_producer
std::unique_ptr< RdKafka::Producer > m_producer
Definition
OpMonPublisher.hpp:79
dunedaq::kafkaopmon::OpMonPublisher::~OpMonPublisher
~OpMonPublisher()
Definition
OpMonPublisher.cpp:56
dunedaq::kafkaopmon::OpMonPublisher::extract_key
std::string extract_key(const dunedaq::opmon::OpMonEntry &e) const noexcept
Definition
OpMonPublisher.hpp:73
dunedaq::kafkaopmon::OpMonPublisher::extract_topic
std::string extract_topic(const dunedaq::opmon::OpMonEntry &) const noexcept
Definition
OpMonPublisher.hpp:72
dunedaq::kafkaopmon::OpMonPublisher::publish
void publish(dunedaq::opmon::OpMonEntry &&) const
Definition
OpMonPublisher.cpp:68
dunedaq::kafkaopmon::OpMonPublisher::OpMonPublisher
OpMonPublisher()=delete
dunedaq::kafkaopmon::OpMonPublisher::m_default_topic
std::string m_default_topic
Definition
OpMonPublisher.hpp:80
dunedaq::opmon::OpMonEntry
Definition
opmon_entry.pb.h:693
dunedaq::kafkaopmon
Definition
OpMonPublisher.hpp:53
dunedaq::FailedConfiguration
FailedConfiguration
Definition
OpMonPublisher.hpp:32
ers::warning
void warning(const Issue &issue)
Definition
ers.hpp:150
ers::error
void error(const Issue &issue)
Definition
ers.hpp:101
Generated on
for DUNE-DAQ by
1.18.0