diff --git a/src/integrations/kafka/consumer.cpp b/src/integrations/kafka/consumer.cpp index 051e0e4b8..48a9ab05d 100644 --- a/src/integrations/kafka/consumer.cpp +++ b/src/integrations/kafka/consumer.cpp @@ -30,6 +30,31 @@ namespace integrations::kafka { namespace { +// TODO: Remove after debugging +std::string time_in_HH_MM_SS_MMM() { + using namespace std::chrono; + + // get current time + auto now = system_clock::now(); + + // get number of milliseconds for the current second + // (remainder after division into seconds) + auto ms = duration_cast(now.time_since_epoch()) % 1000; + + // convert to std::time_t in order to convert to std::tm (broken time) + auto timer = system_clock::to_time_t(now); + + // convert to broken time + std::tm bt = *std::localtime(&timer); + + std::ostringstream oss; + + oss << std::put_time(&bt, "%H:%M:%S"); // HH:MM:SS + oss << '.' << std::setfill('0') << std::setw(3) << ms.count(); + + return oss.str(); +} + utils::BasicResult> GetBatch(RdKafka::KafkaConsumer &consumer, const ConsumerInfo &info, std::atomic &is_running) { @@ -42,7 +67,9 @@ utils::BasicResult> GetBatch(RdKafka::KafkaCon bool run_batch = true; for (int64_t i = 0; remaining_timeout_in_ms > 0 && i < info.batch_size && is_running.load(); ++i) { + spdlog::info("Consuming stream {} ...", time_in_HH_MM_SS_MMM()); std::unique_ptr msg(consumer.consume(remaining_timeout_in_ms)); + spdlog::info("Stream consumed {}", time_in_HH_MM_SS_MMM()); switch (msg->err()) { case RdKafka::ERR__TIMED_OUT: run_batch = false;