From 54771ba1f04c9c3cc816f1d3f508867403032d64 Mon Sep 17 00:00:00 2001 From: yswdqz Date: Wed, 19 Aug 2026 14:34:55 +0800 Subject: [PATCH 1/5] [SDK] BatchSpanProcessor: wait for a full batch instead of draining partial batches - Change the worker wait predicate from '!buffer_.empty()' to 'buffer_.size() >= max_export_batch_size_'. - Export() now only drains the entire buffer when a force flush is pending or the processor is shutting down; on normal wakeups it exports at most one batch of max_export_batch_size spans. This prevents the processor from waking up and draining partial trailing batches every time a span arrives, reducing CPU usage and gRPC request count while preserving ForceFlush/Shutdown drain semantics. --- sdk/src/trace/batch_span_processor.cc | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/sdk/src/trace/batch_span_processor.cc b/sdk/src/trace/batch_span_processor.cc index 403fcb55a..2231b7b98 100644 --- a/sdk/src/trace/batch_span_processor.cc +++ b/sdk/src/trace/batch_span_processor.cc @@ -194,7 +194,7 @@ void BatchSpanProcessor::DoBackgroundWork() // Since `Export()` calls `NotifyCompletion()` which takes `force_flush_cv_m`, // holding `cv_m` while calling `Export()` can lead to a ABBA deadlock. { - // Wait for `timeout` milliseconds. + // Wait for `timeout` milliseconds, or until a full batch is available. std::unique_lock lk(synchronization_data_->cv_m); synchronization_data_->cv.wait_for(lk, timeout, [this] { if (synchronization_data_->is_force_wakeup_background_worker.load( @@ -203,7 +203,7 @@ void BatchSpanProcessor::DoBackgroundWork() return true; } - return !buffer_.empty(); + return buffer_.size() >= max_export_batch_size_; }); synchronization_data_->is_force_wakeup_background_worker.store(false, std::memory_order_release); @@ -248,13 +248,16 @@ void BatchSpanProcessor::Export() } #endif /* ENABLE_THREAD_INSTRUMENTATION_PREVIEW */ + std::uint64_t notify_force_flush = + synchronization_data_->force_flush_pending_sequence.load(std::memory_order_acquire); + bool should_drain = notify_force_flush != 0 || + synchronization_data_->is_shutdown.load(std::memory_order_acquire); + do { std::vector> spans_arr; size_t num_records_to_export{}; - std::uint64_t notify_force_flush = - synchronization_data_->force_flush_pending_sequence.load(std::memory_order_acquire); - if (notify_force_flush) + if (should_drain) { num_records_to_export = buffer_.size(); } @@ -285,7 +288,7 @@ void BatchSpanProcessor::Export() exporter_->Export(nostd::span>(spans_arr.data(), spans_arr.size())); NotifyCompletion(notify_force_flush, exporter_, synchronization_data_); - } while (true); + } while (should_drain); #ifdef ENABLE_THREAD_INSTRUMENTATION_PREVIEW if (worker_thread_instrumentation_ != nullptr) From d003137db961e6ce3a557a553ba0ddd81362b80e Mon Sep 17 00:00:00 2001 From: yswdqz Date: Fri, 21 Aug 2026 14:29:30 +0800 Subject: [PATCH 2/5] Add CHANGELOG entry for BatchSpanProcessor strict batching --- CHANGELOG.md | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4d16d636e..056a8e171 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,6 +15,12 @@ Increment the: ## [Unreleased] +* [SDK] `BatchSpanProcessor` now waits for a full batch + (`max_export_batch_size`) before exporting, instead of draining the buffer + whenever it is non-empty. This reduces gRPC request count and CPU usage + under steady load while preserving `ForceFlush`/`Shutdown` drain semantics. + [#PR_NUMBER](https://github.com/open-telemetry/opentelemetry-cpp/pull/PR_NUMBER) + * [CONFIGURATION] Build the configured resource detectors in SdkBuilder, apply the `detection.attributes` include/exclude filter to the detected attributes, and merge the resource per the resource SDK specification. From 3d02bbf1351c47540f8568ba5ad513e9e6e7263c Mon Sep 17 00:00:00 2001 From: yswdqz Date: Fri, 21 Aug 2026 14:37:28 +0800 Subject: [PATCH 3/5] Update CHANGELOG with PR number 4466 --- CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 056a8e171..0bb646a86 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,7 +19,7 @@ Increment the: (`max_export_batch_size`) before exporting, instead of draining the buffer whenever it is non-empty. This reduces gRPC request count and CPU usage under steady load while preserving `ForceFlush`/`Shutdown` drain semantics. - [#PR_NUMBER](https://github.com/open-telemetry/opentelemetry-cpp/pull/PR_NUMBER) + [#4466](https://github.com/open-telemetry/opentelemetry-cpp/pull/4466) * [CONFIGURATION] Build the configured resource detectors in SdkBuilder, apply the `detection.attributes` include/exclude filter to the detected attributes, From 0e7cbe4f369f8571a2330f748664d4322eab2830 Mon Sep 17 00:00:00 2001 From: yswdqz Date: Sun, 23 Aug 2026 05:21:10 +0800 Subject: [PATCH 4/5] [SDK] BatchSpanProcessor: fix ForceFlush drain logic and respect max_export_batch_size --- sdk/src/trace/batch_span_processor.cc | 18 +- sdk/test/trace/batch_span_processor_test.cc | 177 ++++++++++++++++++++ 2 files changed, 183 insertions(+), 12 deletions(-) diff --git a/sdk/src/trace/batch_span_processor.cc b/sdk/src/trace/batch_span_processor.cc index 2231b7b98..7b502aaf8 100644 --- a/sdk/src/trace/batch_span_processor.cc +++ b/sdk/src/trace/batch_span_processor.cc @@ -250,22 +250,16 @@ void BatchSpanProcessor::Export() std::uint64_t notify_force_flush = synchronization_data_->force_flush_pending_sequence.load(std::memory_order_acquire); - bool should_drain = notify_force_flush != 0 || - synchronization_data_->is_shutdown.load(std::memory_order_acquire); + bool should_drain = + notify_force_flush > + synchronization_data_->force_flush_notified_sequence.load(std::memory_order_acquire) || + synchronization_data_->is_shutdown.load(std::memory_order_acquire); do { std::vector> spans_arr; - size_t num_records_to_export{}; - if (should_drain) - { - num_records_to_export = buffer_.size(); - } - else - { - num_records_to_export = - buffer_.size() >= max_export_batch_size_ ? max_export_batch_size_ : buffer_.size(); - } + size_t num_records_to_export = + buffer_.size() >= max_export_batch_size_ ? max_export_batch_size_ : buffer_.size(); if (num_records_to_export == 0) { diff --git a/sdk/test/trace/batch_span_processor_test.cc b/sdk/test/trace/batch_span_processor_test.cc index 1de105d90..3f82bd35e 100644 --- a/sdk/test/trace/batch_span_processor_test.cc +++ b/sdk/test/trace/batch_span_processor_test.cc @@ -5,8 +5,10 @@ #include #include #include +#include #include #include +#include #include #include #include @@ -105,6 +107,82 @@ class MockSpanExporter final : public sdk::trace::SpanExporter std::chrono::milliseconds export_delay_; }; +class BlockingMockSpanExporter final : public sdk::trace::SpanExporter +{ +public: + BlockingMockSpanExporter( + std::shared_ptr> batch_sizes, + std::shared_ptr> spans_received_count, + std::shared_ptr> is_shutdown, + std::shared_ptr> force_flush_counter = + std::shared_ptr>(new std::atomic(0))) noexcept + : batch_sizes_(std::move(batch_sizes)), + spans_received_count_(std::move(spans_received_count)), + is_shutdown_(std::move(is_shutdown)), + force_flush_counter_(std::move(force_flush_counter)) + {} + + std::unique_ptr MakeRecordable() noexcept override + { + return std::unique_ptr(new sdk::trace::SpanData); + } + + sdk::common::ExportResult Export( + const nostd::span> &recordables) noexcept override + { + { + std::lock_guard lock(mutex_); + batch_sizes_->push_back(recordables.size()); + *spans_received_count_ += recordables.size(); + } + + std::unique_lock lock(mutex_); + cv_.wait(lock, [this] { return !block_export_.load(); }); + lock.unlock(); + + for (auto &recordable : recordables) + { + recordable.release(); + } + + return sdk::common::ExportResult::kSuccess; + } + + bool ForceFlush(std::chrono::microseconds /*timeout*/) noexcept override + { + ++(*force_flush_counter_); + return true; + } + + bool Shutdown(std::chrono::microseconds /* timeout */) noexcept override + { + *is_shutdown_ = true; + return true; + } + + void SetBlock(bool block) + { + { + std::lock_guard lock(mutex_); + block_export_.store(block); + } + if (!block) + { + cv_.notify_all(); + } + } + +private: + std::shared_ptr> batch_sizes_; + std::shared_ptr> spans_received_count_; + std::shared_ptr> is_shutdown_; + std::shared_ptr> force_flush_counter_; + + mutable std::mutex mutex_; + std::condition_variable cv_; + std::atomic block_export_{false}; +}; + /** * Fixture Class */ @@ -165,6 +243,40 @@ TEST_F(BatchSpanProcessorTestPeer, TestShutdown) EXPECT_TRUE(is_shutdown->load()); } +TEST_F(BatchSpanProcessorTestPeer, TestShutdownRespectsMaxExportBatchSize) +{ + std::shared_ptr> batch_sizes(new std::vector()); + std::shared_ptr> spans_received_count(new std::atomic(0)); + std::shared_ptr> is_shutdown(new std::atomic(false)); + std::shared_ptr> force_flush_counter(new std::atomic(0)); + + auto exporter_raw = new BlockingMockSpanExporter(batch_sizes, spans_received_count, is_shutdown, + force_flush_counter); + auto exporter = std::unique_ptr(exporter_raw); + + sdk::trace::BatchSpanProcessorOptions options{}; + options.max_export_batch_size = 100; + options.max_queue_size = 1000; + + auto batch_processor = std::shared_ptr( + new sdk::trace::BatchSpanProcessor(std::move(exporter), options)); + + const int num_spans = 250; + auto test_spans = GetTestSpans(batch_processor, num_spans); + for (int i = 0; i < num_spans; ++i) + { + batch_processor->OnEnd(std::move(test_spans->at(i))); + } + + EXPECT_TRUE(batch_processor->Shutdown()); + + EXPECT_EQ(num_spans, spans_received_count->load()); + for (std::size_t size : *batch_sizes) + { + EXPECT_LE(size, options.max_export_batch_size); + } +} + TEST_F(BatchSpanProcessorTestPeer, TestForceFlush) { std::shared_ptr> shut_down_counter(new std::atomic(0)); @@ -220,6 +332,71 @@ TEST_F(BatchSpanProcessorTestPeer, TestForceFlush) } } +TEST_F(BatchSpanProcessorTestPeer, TestForceFlushDoesNotPermanentlyDrain) +{ + std::shared_ptr> batch_sizes(new std::vector()); + std::shared_ptr> spans_received_count(new std::atomic(0)); + std::shared_ptr> is_shutdown(new std::atomic(false)); + std::shared_ptr> force_flush_counter(new std::atomic(0)); + + auto exporter_raw = new BlockingMockSpanExporter(batch_sizes, spans_received_count, is_shutdown, + force_flush_counter); + auto exporter = std::unique_ptr(exporter_raw); + + sdk::trace::BatchSpanProcessorOptions options{}; + options.max_export_batch_size = 100; + options.schedule_delay_millis = std::chrono::milliseconds(200); + options.max_queue_size = 1000; + + auto batch_processor = std::shared_ptr( + new sdk::trace::BatchSpanProcessor(std::move(exporter), options)); + + auto initial_spans = GetTestSpans(batch_processor, 50); + for (int i = 0; i < 50; ++i) + { + batch_processor->OnEnd(std::move(initial_spans->at(i))); + } + EXPECT_TRUE(batch_processor->ForceFlush()); + EXPECT_EQ(50u, spans_received_count->load()); + + exporter_raw->SetBlock(true); + + auto first_wave = GetTestSpans(batch_processor, 100); + for (int i = 0; i < 100; ++i) + { + batch_processor->OnEnd(std::move(first_wave->at(i))); + } + + auto wait_start = std::chrono::steady_clock::now(); + while (batch_sizes->size() < 2 && + std::chrono::steady_clock::now() - wait_start < std::chrono::seconds(2)) + { + std::this_thread::sleep_for(std::chrono::milliseconds(1)); + } + ASSERT_GE(batch_sizes->size(), 2u); + EXPECT_EQ(100u, batch_sizes->at(1)); + + auto second_wave = GetTestSpans(batch_processor, 5); + for (int i = 0; i < 5; ++i) + { + batch_processor->OnEnd(std::move(second_wave->at(i))); + } + + exporter_raw->SetBlock(false); + std::this_thread::sleep_for(std::chrono::milliseconds(50)); + + EXPECT_EQ(150u, spans_received_count->load()); + EXPECT_EQ(2u, batch_sizes->size()); + + EXPECT_TRUE(batch_processor->ForceFlush()); + EXPECT_EQ(155u, spans_received_count->load()); + + for (std::size_t size : *batch_sizes) + { + EXPECT_LE(size, options.max_export_batch_size); + } +} + // A mock log handler to check whether log messages with a specific level were emitted. struct MockLogHandler : public sdk::common::internal_log::LogHandler { From 005d1adce0253fa3bfecd45dda45a504474211a5 Mon Sep 17 00:00:00 2001 From: yswdqz Date: Sun, 23 Aug 2026 23:51:27 +0800 Subject: [PATCH 5/5] [SDK] BatchSpanProcessor: fix memory leak and data race in test exporter --- sdk/test/trace/batch_span_processor_test.cc | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/sdk/test/trace/batch_span_processor_test.cc b/sdk/test/trace/batch_span_processor_test.cc index 3f82bd35e..f1b01b8f5 100644 --- a/sdk/test/trace/batch_span_processor_test.cc +++ b/sdk/test/trace/batch_span_processor_test.cc @@ -115,11 +115,14 @@ class BlockingMockSpanExporter final : public sdk::trace::SpanExporter std::shared_ptr> spans_received_count, std::shared_ptr> is_shutdown, std::shared_ptr> force_flush_counter = + std::shared_ptr>(new std::atomic(0)), + std::shared_ptr> export_call_count = std::shared_ptr>(new std::atomic(0))) noexcept : batch_sizes_(std::move(batch_sizes)), spans_received_count_(std::move(spans_received_count)), is_shutdown_(std::move(is_shutdown)), - force_flush_counter_(std::move(force_flush_counter)) + force_flush_counter_(std::move(force_flush_counter)), + export_call_count_(std::move(export_call_count)) {} std::unique_ptr MakeRecordable() noexcept override @@ -134,6 +137,7 @@ class BlockingMockSpanExporter final : public sdk::trace::SpanExporter std::lock_guard lock(mutex_); batch_sizes_->push_back(recordables.size()); *spans_received_count_ += recordables.size(); + ++(*export_call_count_); } std::unique_lock lock(mutex_); @@ -142,7 +146,7 @@ class BlockingMockSpanExporter final : public sdk::trace::SpanExporter for (auto &recordable : recordables) { - recordable.release(); + recordable.reset(); } return sdk::common::ExportResult::kSuccess; @@ -177,6 +181,7 @@ class BlockingMockSpanExporter final : public sdk::trace::SpanExporter std::shared_ptr> spans_received_count_; std::shared_ptr> is_shutdown_; std::shared_ptr> force_flush_counter_; + std::shared_ptr> export_call_count_; mutable std::mutex mutex_; std::condition_variable cv_; @@ -338,9 +343,10 @@ TEST_F(BatchSpanProcessorTestPeer, TestForceFlushDoesNotPermanentlyDrain) std::shared_ptr> spans_received_count(new std::atomic(0)); std::shared_ptr> is_shutdown(new std::atomic(false)); std::shared_ptr> force_flush_counter(new std::atomic(0)); + std::shared_ptr> export_call_count(new std::atomic(0)); auto exporter_raw = new BlockingMockSpanExporter(batch_sizes, spans_received_count, is_shutdown, - force_flush_counter); + force_flush_counter, export_call_count); auto exporter = std::unique_ptr(exporter_raw); sdk::trace::BatchSpanProcessorOptions options{}; @@ -368,12 +374,12 @@ TEST_F(BatchSpanProcessorTestPeer, TestForceFlushDoesNotPermanentlyDrain) } auto wait_start = std::chrono::steady_clock::now(); - while (batch_sizes->size() < 2 && + while (export_call_count->load() < 2 && std::chrono::steady_clock::now() - wait_start < std::chrono::seconds(2)) { std::this_thread::sleep_for(std::chrono::milliseconds(1)); } - ASSERT_GE(batch_sizes->size(), 2u); + ASSERT_GE(export_call_count->load(), 2u); EXPECT_EQ(100u, batch_sizes->at(1)); auto second_wave = GetTestSpans(batch_processor, 5);