DUNE-DAQ
DUNE Trigger and Data Acquisition software
Toggle main menu visibility
Loading...
Searching...
No Matches
dunedaq
sourcecode
erskafka
src
ERSPublisher.cpp
Go to the documentation of this file.
1
8
9
#include "
erskafka/ERSPublisher.hpp
"
10
11
#include <iostream>
12
13
#include <stdexcept>
14
15
using namespace
dunedaq::erskafka
;
16
17
ERSPublisher::ERSPublisher
(
const
nlohmann::json& conf)
18
{
19
20
RdKafka::Conf* k_conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
21
std::string errstr;
22
23
auto
it = conf.find(
"bootstrap"
);
24
if
(it == conf.end()) {
25
std::cerr <<
"Missing bootstrap from json file"
;
26
throw
std::runtime_error(
"Missing bootstrap from json file"
);
27
}
28
29
k_conf->set(
"bootstrap.servers"
, *it, errstr);
30
if
(errstr !=
""
) {
31
throw
std::runtime_error(errstr);
32
}
33
34
std::string client_id;
35
it = conf.find(
"cliend_id"
);
36
if
(it != conf.end())
37
client_id = *it;
38
else
if
(
const
char
* env_p = std::getenv(
"DUNEDAQ_APPLICATION_NAME"
))
39
client_id = env_p;
40
else
41
client_id =
"erskafkaproducerdefault"
;
42
43
k_conf->set(
"client.id"
, client_id, errstr);
44
if
(errstr !=
""
) {
45
throw
std::runtime_error(errstr);
46
}
47
48
// Create producer instance
49
m_producer
.reset(RdKafka::Producer::create(k_conf, errstr));
50
51
if
(errstr !=
""
) {
52
throw
std::runtime_error(errstr);
53
}
54
55
it = conf.find(
"default_topic"
);
56
if
(it != conf.end())
57
m_default_topic
= *it;
58
}
59
60
bool
61
ERSPublisher::publish
(ersschema::IssueChain&& issue)
const
62
{
63
64
std::string binary;
65
issue.SerializeToString(&binary);
66
67
// get the topic
68
auto
topic
=
ERSPublisher::topic
(issue);
69
70
auto
key
=
ERSPublisher::key
(issue);
71
72
// RdKafka::Producer::RK_MSG_COPY to be investigated
73
RdKafka::ErrorCode
err
=
m_producer
->produce(
topic
,
74
RdKafka::Topic::PARTITION_UA,
75
RdKafka::Producer::RK_MSG_COPY,
76
const_cast<
char
*
>
(binary.c_str()),
77
binary.size(),
78
key
.c_str(),
79
key
.size(),
80
0,
81
nullptr
);
82
if
(err != RdKafka::ERR_NO_ERROR) {
83
return
false
;
84
}
85
86
return
true
;
87
}
ERSPublisher.hpp
ERSPublisher.ERSPublisher
Definition
ERSPublisher.py:19
dunedaq::erskafka::ERSPublisher::publish
bool publish(dunedaq::ersschema::IssueChain &&) const
dunedaq::erskafka::ERSPublisher::m_producer
std::unique_ptr< RdKafka::Producer > m_producer
Definition
ERSPublisher.hpp:51
dunedaq::erskafka::ERSPublisher::m_default_topic
std::string m_default_topic
Definition
ERSPublisher.hpp:52
dunedaq::erskafka::ERSPublisher::key
std::string key(const dunedaq::ersschema::IssueChain &i) const
Definition
ERSPublisher.hpp:48
dunedaq::erskafka::ERSPublisher::topic
std::string topic(const dunedaq::ersschema::IssueChain &) const
Definition
ERSPublisher.hpp:46
dunedaq::erskafka
Definition
ERSPublisher.hpp:25
fanoutRealLinks.err
err
Definition
fanoutRealLinks.py:10
Generated on
for DUNE-DAQ by
1.18.0