Skip to content

[SDK] BatchSpanProcessor: wait for a full batch instead of draining partial batches - #4466

Open
yswdqz wants to merge 9 commits into
open-telemetry:mainfrom
yswdqz:fix/bsp-strict-batch
Open

[SDK] BatchSpanProcessor: wait for a full batch instead of draining partial batches#4466
yswdqz wants to merge 9 commits into
open-telemetry:mainfrom
yswdqz:fix/bsp-strict-batch

Conversation

@yswdqz

@yswdqz yswdqz commented Aug 21, 2026

Copy link
Copy Markdown

fix #4449

Problem

BatchSpanProcessor currently wakes up whenever the buffer is non-empty and
then drains the entire buffer in a tight loop. Under steady load this produces
exports that are much smaller than max_export_batch_size, causing:

  • Up to 2x more gRPC export requests than necessary
  • Higher CPU usage due to repeated serialization and wakeups

Changes

  • 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 preserves ForceFlush/Shutdown semantics while making the normal
export path strictly batch-oriented.

Performance

Metric Official Custom
Achieved SPS 19999.66 19999.99
Spans received 1,200,000 1,200,000
Export requests 254,683 586
Avg batch size 4.7 2047.8
Mid‑15‑s CPU usage 0.512 cores 0.200 cores

I constructed a test where one span is finished every 50 µs, and the above metrics capture the performance difference between the two implementations under that steady load. In real-world scenarios, spans tend to be completed in batches, so the performance gap would be much less pronounced under typical conditions.

Checklist

  • CHANGELOG.md updated for non-trivial changes
  • Unit tests have been added (No new tests have been added, existing batch_span_processor_test and
    batch_span_processor_test_stress all pass)
  • Changes in public API reviewed (no public API changes;)

@yswdqz
yswdqz requested a review from a team as a code owner August 21, 2026 06:36
@linux-foundation-easycla

linux-foundation-easycla Bot commented Aug 21, 2026

Copy link
Copy Markdown

CLA Signed
The committers listed above are authorized under a signed CLA.

yswdqz added 3 commits August 21, 2026 14:41
…artial 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.
@yswdqz
yswdqz force-pushed the fix/bsp-strict-batch branch from c49d86c to 3d02bbf Compare August 21, 2026 06:41
@codecov

codecov Bot commented Aug 21, 2026

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 83.11%. Comparing base (93f16f3) to head (f048172).
⚠️ Report is 11 commits behind head on main.

Additional details and impacted files

Impacted file tree graph

@@            Coverage Diff             @@
##             main    #4466      +/-   ##
==========================================
+ Coverage   82.62%   83.11%   +0.50%     
==========================================
  Files         512      519       +7     
  Lines       20138    20254     +116     
==========================================
+ Hits        16637    16833     +196     
+ Misses       3501     3421      -80     
Files with missing lines Coverage Δ
sdk/src/trace/batch_span_processor.cc 85.32% <100.00%> (-0.16%) ⬇️

... and 21 files with indirect coverage changes

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

@marcalff

marcalff commented Aug 21, 2026

Copy link
Copy Markdown
Member

Thanks for the fix.

Please see clang-format errors, either run clang-format or apply this manually:

diff --git a/sdk/src/trace/batch_span_processor.cc b/sdk/src/trace/batch_span_processor.cc
index 2231b7b..1a1a7d8 100644
--- a/sdk/src/trace/batch_span_processor.cc
+++ b/sdk/src/trace/batch_span_processor.cc
@@ -250,8 +250,8 @@ 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 != 0 || synchronization_data_->is_shutdown.load(std::memory_order_acquire);
 
   do
   {

Comment thread sdk/src/trace/batch_span_processor.cc Outdated
Comment on lines +253 to +254
bool should_drain = notify_force_flush != 0 ||
synchronization_data_->is_shutdown.load(std::memory_order_acquire);

@denizariyan denizariyan Aug 21, 2026

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

notify_force_flush (i.e., synchronization_data_->force_flush_pending_sequence) is a monotonically increasing counter that increases on every call to ForceFlush. This would mean this condition would be permanently true after the first time one calls BatchSpanProcessor::ForceFlush.

Hence while (should_drain) below never becomes while (false), so Export() keeps draining to empty on every call which means the one-batch-per-wakeup behavior this PR adds never takes effect after the first call to BatchSpanProcessor::ForceFlush.

I think something like this would actually solve that problem.

Suggested change
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);

Maybe you could add some tests to confirm/check?

This issue could be tested with a test exporter that records how many spans each Export() call receives and pauses inside the first call, so the worker is held mid-export and you control what's in the buffer when it resumes. You could then add a couple spans while it's paused mid export to see that it will still drain less than the max batch size immediately instead of returning after exporting one batch as the PR suggests IF the BatchSpanProcessor::ForceFlush was ever called.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch, thank you for the review~ Updated the condition to compare against force_flush_notified_sequence and added a regression test covering the scenario you described.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks!

Comment thread sdk/src/trace/batch_span_processor.cc Outdated
std::uint64_t notify_force_flush =
synchronization_data_->force_flush_pending_sequence.load(std::memory_order_acquire);
if (notify_force_flush)
if (should_drain)

@ThomsonTan ThomsonTan Aug 21, 2026

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

should_drain is doing two orthogonal jobs here: how many records to take, and whether to keep looping. Because is_shutdown alone now sets it, it hands the entire buffer to a single Export() call, bypassing max_export_batch_size. Before this PR, a program that never called ForceFlush had notify_force_flush == 0, so DrainQueue() still chunked; now a backlog of 3000 with max_export_batch_size = 2048 goes out as one 3000-span request instead of 2048 + 952. If the receiver rejects the oversized message (max_receive_message_length defaults to 4 MiB), the spans are already consumed and the result is discarded at the line - and at process exit there is no retry.

Suggest keeping the cap unconditional and letting should_drain control only the loop.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You're right, thank you for the review~ The cap is now unconditional and should_drain only controls the loop. Added a separate test for shutdown draining as well.

@yswdqz
yswdqz force-pushed the fix/bsp-strict-batch branch from adb1e93 to 0e7cbe4 Compare August 22, 2026 21:26
@yswdqz

yswdqz commented Aug 22, 2026

Copy link
Copy Markdown
Author

Apologies for my oversight and insufficient testing. Thanks for the reviews. I've addressed all the feedback: fixed the clang-format issue, corrected the ForceFlush condition, kept the batch size cap unconditional, and added regression tests for both issues.

@yswdqz

yswdqz commented Aug 23, 2026

Copy link
Copy Markdown
Author

Some CI jobs failed because the new BlockingMockSpanExporter had a memory leak and a data race. Pushed a fix and verified locally with ASan and TSan — all batch_span_processor_test tests pass now.

Comment thread sdk/src/trace/batch_span_processor.cc Outdated
@@ -285,7 +282,7 @@ void BatchSpanProcessor::Export()

exporter_->Export(nostd::span<std::unique_ptr<Recordable>>(spans_arr.data(), spans_arr.size()));
NotifyCompletion(notify_force_flush, exporter_, synchronization_data_);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Looking a bit deeper I think we also have an issue here in regards to handling of ForceFlush.

When more than max_export_batch_size_ spans are queued at the time, ForceFlush will return after exporting the first batch, hence returning before all the spans are exported because the notify is called within the loop.

To be clear, main can also return from ForceFlush() with spans in the buffer, but that's only for ones that arrived during the export, which is fine IMO. What's new here is spans that were already buffered when ForceFlush() was called can be being left behind.

Because the spec only mentions handling tasks (i.e., spans) received prior to the call to the ForceFlush I think this calls for a "snapshot-then-export" like pattern.

Maybe something like below, WDYT?

Export():
    notify_force_flush = pending_sequence.load()
    should_drain       = notify_force_flush > notified_sequence.load() || is_shutdown.load()

    # snapshot the target ONCE, before exporting anything
    remaining = should_drain ? buffer_.size()
                             : min(buffer_.size(), max_export_batch_size_)

    while remaining > 0:
        n = min(remaining, max_export_batch_size_)   # cap still applies on every path
        consume n from buffer_ into spans_arr
        exporter_->Export(spans_arr)
        remaining -= n                                # NOT re-reading buffer_.size()

    NotifyCompletion(notify_force_flush, ...)          # once, after the snapshot is out

The snapshot pattern would solve a couple problems:

  • A clause that loops until empty where we keep checking the buffer size within the loop races against the producer so under high sustained load it may never get to finish exporting (like a livelock situation), eventually running into the timeout. Snapshotting at the call time matches both what the spec defines and prevents this potential livelock.
  • Moving the notify to the end of the loop covers all cases without branching like we have today.
  • And IMO it is easier to read in general.

Maybe we could also add some tests to ensure this doesn't regress later?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@denizariyan You're right, that's a real issue — the in-loop notify lets ForceFlush return after the first batch when multiple batches are queued. Implemented your snapshot-then-export proposal exactly (snapshot once, remaining -= n, notify once after the loop), which also bounds the work per call. Added TestForceFlushExportsAllBufferedSpans (250 spans / batch 100): verified it fails on the old code (returns at 200 received) and passes with the fix. ASAN/TSAN/UBSAN green locally.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@denizariyan Sorry for the back-and-forth. I can't run the full CI locally, so some issues only surface once CI runs and I can only fix them after seeing the failures. I'll run everything I can locally (format, ASAN/TSAN/UBSAN, relevant unit tests) before pushing from now on. Thanks for your patience!

@denizariyan denizariyan left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM, thanks.

We have a similar pattern in BatchLogRecordProcessor::Export() too which is now different compared to the trace path, would be nice to create an issue so we can close the deviation on it in a follow up

@yswdqz

yswdqz commented Aug 28, 2026

Copy link
Copy Markdown
Author

@denizariyan You are right, I confirmed locally that BatchLogRecordProcessor::Export() has the same tight-loop drain behavior, so it is now inconsistent with the trace path. I filed #4498 to track this.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[SDK] BatchSpanProcessor drain the queue in a tight loop instead of waiting for the next full batch

4 participants