mirror of
https://github.com/Telecominfraproject/wlan-cloud-lib-cppkafka.git
synced 2025-11-01 02:57:53 +00:00
104 lines
3.0 KiB
C++
104 lines
3.0 KiB
C++
#include <stdexcept>
|
|
#include <iostream>
|
|
#include <csignal>
|
|
#include <boost/program_options.hpp>
|
|
#include "cppkafka/consumer.h"
|
|
#include "cppkafka/configuration.h"
|
|
|
|
using std::string;
|
|
using std::exception;
|
|
using std::cout;
|
|
using std::endl;
|
|
|
|
using cppkafka::Consumer;
|
|
using cppkafka::Configuration;
|
|
using cppkafka::Message;
|
|
using cppkafka::TopicPartitionList;
|
|
|
|
namespace po = boost::program_options;
|
|
|
|
bool running = true;
|
|
|
|
int main(int argc, char* argv[]) {
|
|
string brokers;
|
|
string topic_name;
|
|
string group_id;
|
|
|
|
po::options_description options("Options");
|
|
options.add_options()
|
|
("help,h", "produce this help message")
|
|
("brokers", po::value<string>(&brokers)->required(),
|
|
"the kafka broker list")
|
|
("topic", po::value<string>(&topic_name)->required(),
|
|
"the topic in which to write to")
|
|
("group-id", po::value<string>(&group_id)->required(),
|
|
"the consumer group id")
|
|
;
|
|
|
|
po::variables_map vm;
|
|
|
|
try {
|
|
po::store(po::command_line_parser(argc, argv).options(options).run(), vm);
|
|
po::notify(vm);
|
|
}
|
|
catch (exception& ex) {
|
|
cout << "Error parsing options: " << ex.what() << endl;
|
|
cout << endl;
|
|
cout << options << endl;
|
|
return 1;
|
|
}
|
|
|
|
// Stop processing on SIGINT
|
|
signal(SIGINT, [](int) { running = false; });
|
|
|
|
// Construct the configuration
|
|
Configuration config = {
|
|
{ "metadata.broker.list", brokers },
|
|
{ "group.id", group_id },
|
|
// Disable auto commit
|
|
{ "enable.auto.commit", false }
|
|
};
|
|
|
|
// Create the consumer
|
|
Consumer consumer(config);
|
|
|
|
// Print the assigned partitions on assignment
|
|
consumer.set_assignment_callback([](const TopicPartitionList& partitions) {
|
|
cout << "Got assigned: " << partitions << endl;
|
|
});
|
|
|
|
// Print the revoked partitions on revocation
|
|
consumer.set_revocation_callback([](const TopicPartitionList& partitions) {
|
|
cout << "Got revoked: " << partitions << endl;
|
|
});
|
|
|
|
// Subscribe to the topic
|
|
consumer.subscribe({ topic_name });
|
|
|
|
cout << "Consuming messages from topic " << topic_name << endl;
|
|
|
|
// Now read lines and write them into kafka
|
|
while (running) {
|
|
// Try to consume a message
|
|
Message msg = consumer.poll();
|
|
if (msg) {
|
|
// If we managed to get a message
|
|
if (msg.has_error()) {
|
|
if (msg.get_error() != RD_KAFKA_RESP_ERR__PARTITION_EOF) {
|
|
cout << "[+] Received error notification: " << msg.get_error_string() << endl;
|
|
}
|
|
}
|
|
else {
|
|
// Print the key (if any)
|
|
if (msg.get_key()) {
|
|
cout << msg.get_key() << " -> ";
|
|
}
|
|
// Print the payload
|
|
cout << msg.get_payload() << endl;
|
|
// Now commit the message
|
|
consumer.commit(msg);
|
|
}
|
|
}
|
|
}
|
|
}
|