mirror of
https://github.com/Telecominfraproject/wlan-cloud-lib-cppkafka.git
synced 2025-11-02 19:47:55 +00:00
Add examples
This commit is contained in:
@@ -50,3 +50,4 @@ add_subdirectory(src)
|
|||||||
add_dependencies(cppkafka googletest)
|
add_dependencies(cppkafka googletest)
|
||||||
enable_testing()
|
enable_testing()
|
||||||
add_subdirectory(tests)
|
add_subdirectory(tests)
|
||||||
|
add_subdirectory(examples)
|
||||||
|
|||||||
8
examples/CMakeLists.txt
Normal file
8
examples/CMakeLists.txt
Normal file
@@ -0,0 +1,8 @@
|
|||||||
|
find_package(Boost REQUIRED COMPONENTS program_options)
|
||||||
|
|
||||||
|
link_libraries(${Boost_LIBRARIES} cppkafka ${RDKAFKA_LIBRARY} ${ZOOKEEPER_LIBRARY})
|
||||||
|
|
||||||
|
include_directories(${CMAKE_CURRENT_SOURCE_DIR}/../include)
|
||||||
|
|
||||||
|
add_executable(kafka_producer kafka_producer.cpp)
|
||||||
|
add_executable(kafka_consumer kafka_consumer.cpp)
|
||||||
93
examples/kafka_consumer.cpp
Normal file
93
examples/kafka_consumer.cpp
Normal file
@@ -0,0 +1,93 @@
|
|||||||
|
#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;
|
||||||
|
|
||||||
|
namespace po = boost::program_options;
|
||||||
|
|
||||||
|
#ifndef CPPKAFKA_HAVE_ZOOKEEPER
|
||||||
|
static_assert(false, "Examples require the zookeeper extension");
|
||||||
|
#endif
|
||||||
|
|
||||||
|
bool running = true;
|
||||||
|
|
||||||
|
int main(int argc, char* argv[]) {
|
||||||
|
string zookeeper_endpoint;
|
||||||
|
string topic_name;
|
||||||
|
string group_id;
|
||||||
|
|
||||||
|
po::options_description options("Options");
|
||||||
|
options.add_options()
|
||||||
|
("help,h", "produce this help message")
|
||||||
|
("zookeeper", po::value<string>(&zookeeper_endpoint)->required(),
|
||||||
|
"the zookeeper endpoint")
|
||||||
|
("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;
|
||||||
|
config.set("zookeeper", zookeeper_endpoint);
|
||||||
|
config.set("group.id", group_id);
|
||||||
|
// Disable auto commit
|
||||||
|
config.set("enable.auto.commit", false);
|
||||||
|
|
||||||
|
// Create the consumer
|
||||||
|
Consumer consumer(config);
|
||||||
|
|
||||||
|
// Subscribe to the topic
|
||||||
|
consumer.subscribe({ topic_name });
|
||||||
|
|
||||||
|
// 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().as_string() << " -> ";
|
||||||
|
}
|
||||||
|
// Print the payload
|
||||||
|
cout << msg.get_payload().as_string() << endl;
|
||||||
|
// Now commit the message
|
||||||
|
consumer.commit(msg);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
76
examples/kafka_producer.cpp
Normal file
76
examples/kafka_producer.cpp
Normal file
@@ -0,0 +1,76 @@
|
|||||||
|
#include <stdexcept>
|
||||||
|
#include <iostream>
|
||||||
|
#include <boost/program_options.hpp>
|
||||||
|
#include "cppkafka/producer.h"
|
||||||
|
#include "cppkafka/configuration.h"
|
||||||
|
|
||||||
|
using std::string;
|
||||||
|
using std::exception;
|
||||||
|
using std::getline;
|
||||||
|
using std::cin;
|
||||||
|
using std::cout;
|
||||||
|
using std::endl;
|
||||||
|
|
||||||
|
using cppkafka::Producer;
|
||||||
|
using cppkafka::Configuration;
|
||||||
|
using cppkafka::Topic;
|
||||||
|
using cppkafka::Partition;
|
||||||
|
|
||||||
|
namespace po = boost::program_options;
|
||||||
|
|
||||||
|
#ifndef CPPKAFKA_HAVE_ZOOKEEPER
|
||||||
|
static_assert(false, "Examples require the zookeeper extension");
|
||||||
|
#endif
|
||||||
|
|
||||||
|
int main(int argc, char* argv[]) {
|
||||||
|
string zookeeper_endpoint;
|
||||||
|
string topic_name;
|
||||||
|
int partition_value = -1;
|
||||||
|
|
||||||
|
po::options_description options("Options");
|
||||||
|
options.add_options()
|
||||||
|
("help,h", "produce this help message")
|
||||||
|
("zookeeper", po::value<string>(&zookeeper_endpoint)->required(),
|
||||||
|
"the zookeeper endpoint")
|
||||||
|
("topic", po::value<string>(&topic_name)->required(),
|
||||||
|
"the topic in which to write to")
|
||||||
|
("partition", po::value<int>(&partition_value),
|
||||||
|
"the partition to write into (unassigned if not provided)")
|
||||||
|
;
|
||||||
|
|
||||||
|
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;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Get the partition we want to write to. If no partition is provided, this will be
|
||||||
|
// an unassigned one
|
||||||
|
Partition partition;
|
||||||
|
if (partition_value != -1) {
|
||||||
|
partition = partition_value;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Construct the configuration
|
||||||
|
Configuration config;
|
||||||
|
config.set("zookeeper", zookeeper_endpoint);
|
||||||
|
|
||||||
|
// Create the producer
|
||||||
|
Producer producer(config);
|
||||||
|
// Get the topic we want
|
||||||
|
Topic topic = producer.get_topic(topic_name);
|
||||||
|
|
||||||
|
// Now read lines and write them into kafka
|
||||||
|
string line;
|
||||||
|
while (getline(cin, line)) {
|
||||||
|
// Write the string into the partition
|
||||||
|
producer.produce(topic, partition, line);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -92,6 +92,11 @@ public:
|
|||||||
*/
|
*/
|
||||||
size_t get_size() const;
|
size_t get_size() const;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Checks whether this is a non empty buffer
|
||||||
|
*/
|
||||||
|
explicit operator bool() const;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Converts the contents of the buffer into a string
|
* Converts the contents of the buffer into a string
|
||||||
*/
|
*/
|
||||||
|
|||||||
@@ -84,6 +84,11 @@ public:
|
|||||||
*/
|
*/
|
||||||
rd_kafka_resp_err_t get_error() const;
|
rd_kafka_resp_err_t get_error() const;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Gets the error as a string
|
||||||
|
*/
|
||||||
|
std::string get_error_string() const;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Gets the topic that this message belongs to
|
* Gets the topic that this message belongs to
|
||||||
*/
|
*/
|
||||||
|
|||||||
@@ -51,6 +51,10 @@ size_t Buffer::get_size() const {
|
|||||||
return size_;
|
return size_;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
Buffer::operator bool() const {
|
||||||
|
return data_ != nullptr;
|
||||||
|
}
|
||||||
|
|
||||||
string Buffer::as_string() const {
|
string Buffer::as_string() const {
|
||||||
return string(data_, data_ + size_);
|
return string(data_, data_ + size_);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -71,6 +71,10 @@ rd_kafka_resp_err_t Message::get_error() const {
|
|||||||
return handle_->err;
|
return handle_->err;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
string Message::get_error_string() const {
|
||||||
|
return rd_kafka_err2str(handle_->err);
|
||||||
|
}
|
||||||
|
|
||||||
int Message::get_partition() const {
|
int Message::get_partition() const {
|
||||||
return handle_->partition;
|
return handle_->partition;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user