DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
KafkaStream.cpp
Go to the documentation of this file.
1
8
9#include "KafkaStream.hpp"
10#include <boost/crc.hpp>
11#include <chrono>
12#include <ers/StreamFactory.hpp>
13#include <iostream>
14#include <string>
15#include <vector>
16
18
19
22namespace erskafka {
23erskafka::KafkaStream::KafkaStream(const std::string& param)
24{
25
26 if (const char* env_p = std::getenv("DUNEDAQ_PARTITION"))
27 m_partition = env_p;
28
29 // Kafka server settings
30 std::string brokers = param;
31 std::string errstr;
32
33 RdKafka::Conf* conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
34 conf->set("bootstrap.servers", brokers, errstr);
35 if (errstr != "") {
36 std::cout << "Bootstrap server error : " << errstr << '\n';
37 }
38 if (const char* env_p = std::getenv("DUNEDAQ_APPLICATION_NAME"))
39 conf->set("client.id", env_p, errstr);
40 else
41 conf->set("client.id", "erskafkaproducerdefault", errstr);
42 if (errstr != "") {
43 std::cout << "Producer configuration error : " << errstr << '\n';
44 }
45 // Create producer instance
46 m_producer = RdKafka::Producer::create(conf, errstr);
47
48 if (errstr != "") {
49 std::cout << "Producer creation error : " << errstr << '\n';
50 }
51}
52
53void
54erskafka::KafkaStream::ers_to_json(const ers::Issue& issue, size_t chain, std::vector<nlohmann::json>& j_objs)
55{
56 try {
57 nlohmann::json message;
58 message["partition"] = m_partition.c_str();
59 message["issue_name"] = issue.get_class_name();
60 message["message"] = issue.message().c_str();
61 message["severity"] = ers::to_string(issue.severity());
62 message["usecs_since_epoch"] =
63 std::chrono::duration_cast<std::chrono::microseconds>(issue.ptime().time_since_epoch()).count();
64 message["time"] = std::chrono::duration_cast<std::chrono::milliseconds>(issue.ptime().time_since_epoch()).count();
65
66 message["qualifiers"] = issue.qualifiers();
67 message["params"] = nlohmann::json::array({});
68 for (auto p : issue.parameters()) {
69 message["params"].push_back(p.first + ": " + p.second);
70 }
71 message["cwd"] = issue.context().cwd();
72 message["file_name"] = issue.context().file_name();
73 message["function_name"] = issue.context().function_name();
74 message["host_name"] = issue.context().host_name();
75 message["package_name"] = issue.context().package_name();
76 message["user_name"] = issue.context().user_name();
77 message["application_name"] = issue.context().application_name();
78 message["user_id"] = issue.context().user_id();
79 message["process_id"] = issue.context().process_id();
80 message["thread_id"] = issue.context().thread_id();
81 message["line_number"] = issue.context().line_number();
82 message["chain"] = chain;
83
84 if (issue.cause()) {
85 ers_to_json(*issue.cause(), 2, j_objs);
86 }
87 j_objs.push_back(message);
88 } catch (const std::exception& e) {
89 std::cout << "Conversion from json error : " << e.what() << '\n';
90 }
91}
92
93void
94erskafka::KafkaStream::kafka_exporter(std::string input, std::string topic)
95{
96 try {
97 // RdKafka::Producer::RK_MSG_COPY to be investigated
98 RdKafka::ErrorCode err = m_producer->produce(topic,
99 RdKafka::Topic::PARTITION_UA,
100 RdKafka::Producer::RK_MSG_COPY,
101 const_cast<char*>(input.c_str()),
102 input.size(),
103 nullptr,
104 0,
105 0,
106 nullptr,
107 nullptr);
108 if (err != RdKafka::ERR_NO_ERROR) {
109 std::cout << "% Failed to produce to topic " << topic << ": " << RdKafka::err2str(err) << std::endl;
110 }
111 } catch (const std::exception& e) {
112 std::cout << "Producer error : " << e.what() << '\n';
113 }
114}
115
119void
121{
122 try {
123 std::vector<nlohmann::json> j_objs;
124 if (issue.cause() == nullptr) {
125 ers_to_json(issue, 0, j_objs);
126 } else {
127 ers_to_json(issue, 1, j_objs);
128 }
129
130 // build a unique hash for a group of nested issues
131 std::ostringstream tmpstream(issue.message());
132 tmpstream << issue.context().process_id() << issue.time_t() << issue.context().application_name()
133 << issue.context().host_name() << rand();
134 std::string tmp = tmpstream.str();
135 boost::crc_32_type crc32;
136 crc32.process_bytes(tmp.c_str(), tmp.length());
137
138 for (auto j : j_objs) {
139 j["group_hash"] = crc32.checksum();
140
141 erskafka::KafkaStream::kafka_exporter(j.dump(), "erskafka-reporting");
142 }
143 chained().write(issue);
144 } catch (const std::exception& e) {
145 std::cout << "Producer error : " << e.what() << '\n';
146 }
147}
148} // namespace erskafka
virtual int line_number() const =0
virtual const char * user_name() const =0
virtual int user_id() const =0
virtual pid_t thread_id() const =0
virtual const char * host_name() const =0
virtual pid_t process_id() const =0
virtual const char * package_name() const =0
virtual const char * application_name() const =0
virtual const char * file_name() const =0
virtual const char * cwd() const =0
virtual const char * function_name() const =0
Base class for any user define issue.
Definition Issue.hpp:76
const Context & context() const
Context of the issue.
Definition Issue.hpp:102
ers::Severity severity() const
severity of the issue
Definition Issue.hpp:110
virtual const char * get_class_name() const =0
Get key for class (used for serialisation).
const system_clock::time_point & ptime() const
original time point of the issue
Definition Issue.hpp:131
const std::vector< std::string > & qualifiers() const
return array of qualifiers
Definition Issue.hpp:106
const std::string & message() const
General cause of the issue.
Definition Issue.hpp:104
const string_map & parameters() const
return array of parameters
Definition Issue.hpp:108
std::time_t time_t() const
seconds since 1 Jan 1970
Definition Issue.cpp:146
const Issue * cause() const
return the cause Issue of this Issue
Definition Issue.hpp:100
OutputStream & chained()
ES stream implementation.
RdKafka::Producer * m_producer
void write(const ers::Issue &issue) override
void kafka_exporter(std::string input, std::string topic)
KafkaStream(const std::string &param)
void ers_to_json(const ers::Issue &issue, size_t chain, std::vector< nlohmann::json > &j_objs)
#define ERS_REGISTER_OUTPUT_STREAM(class, name, param)
Definition macro.hpp:27
std::string to_string(severity s)