Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
[#4466](https://github.com/open-telemetry/opentelemetry-cpp/pull/4466)

* [CONFIGURATION] Add a configuration builder for the host resource detector
[#4451](https://github.com/open-telemetry/opentelemetry-cpp/issues/4451)
* [CONFIGURATION] Build the configured resource detectors in SdkBuilder, apply
Expand Down
27 changes: 12 additions & 15 deletions sdk/src/trace/batch_span_processor.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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<std::mutex> lk(synchronization_data_->cv_m);
synchronization_data_->cv.wait_for(lk, timeout, [this] {
if (synchronization_data_->is_force_wakeup_background_worker.load(
Expand All @@ -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);
Expand Down Expand Up @@ -248,21 +248,18 @@ 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 >
synchronization_data_->force_flush_notified_sequence.load(std::memory_order_acquire) ||
synchronization_data_->is_shutdown.load(std::memory_order_acquire);

do
{
std::vector<std::unique_ptr<Recordable>> 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)
{
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)
{
Expand All @@ -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_);
} while (true);
} while (should_drain);

#ifdef ENABLE_THREAD_INSTRUMENTATION_PREVIEW
if (worker_thread_instrumentation_ != nullptr)
Expand Down
183 changes: 183 additions & 0 deletions sdk/test/trace/batch_span_processor_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,10 @@
#include <algorithm>
#include <atomic>
#include <chrono>
#include <condition_variable>
#include <cstddef>
#include <cstdlib>
#include <mutex>
#include <string>
#include <thread>
#include <utility>
Expand Down Expand Up @@ -105,6 +107,87 @@ class MockSpanExporter final : public sdk::trace::SpanExporter
std::chrono::milliseconds export_delay_;
};

class BlockingMockSpanExporter final : public sdk::trace::SpanExporter
{
public:
BlockingMockSpanExporter(
std::shared_ptr<std::vector<std::size_t>> batch_sizes,
std::shared_ptr<std::atomic<std::size_t>> spans_received_count,
std::shared_ptr<std::atomic<bool>> is_shutdown,
std::shared_ptr<std::atomic<std::size_t>> force_flush_counter =
std::shared_ptr<std::atomic<std::size_t>>(new std::atomic<std::size_t>(0)),
std::shared_ptr<std::atomic<std::size_t>> export_call_count =
std::shared_ptr<std::atomic<std::size_t>>(new std::atomic<std::size_t>(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)),
export_call_count_(std::move(export_call_count))
{}

std::unique_ptr<sdk::trace::Recordable> MakeRecordable() noexcept override
{
return std::unique_ptr<sdk::trace::Recordable>(new sdk::trace::SpanData);
}

sdk::common::ExportResult Export(
const nostd::span<std::unique_ptr<sdk::trace::Recordable>> &recordables) noexcept override
{
{
std::lock_guard<std::mutex> lock(mutex_);
batch_sizes_->push_back(recordables.size());
*spans_received_count_ += recordables.size();
++(*export_call_count_);
}

std::unique_lock<std::mutex> lock(mutex_);
cv_.wait(lock, [this] { return !block_export_.load(); });
lock.unlock();

for (auto &recordable : recordables)
{
recordable.reset();
}

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<std::mutex> lock(mutex_);
block_export_.store(block);
}
if (!block)
{
cv_.notify_all();
}
}

private:
std::shared_ptr<std::vector<std::size_t>> batch_sizes_;
std::shared_ptr<std::atomic<std::size_t>> spans_received_count_;
std::shared_ptr<std::atomic<bool>> is_shutdown_;
std::shared_ptr<std::atomic<std::size_t>> force_flush_counter_;
std::shared_ptr<std::atomic<std::size_t>> export_call_count_;

mutable std::mutex mutex_;
std::condition_variable cv_;
std::atomic<bool> block_export_{false};
};

/**
* Fixture Class
*/
Expand Down Expand Up @@ -165,6 +248,40 @@ TEST_F(BatchSpanProcessorTestPeer, TestShutdown)
EXPECT_TRUE(is_shutdown->load());
}

TEST_F(BatchSpanProcessorTestPeer, TestShutdownRespectsMaxExportBatchSize)
{
std::shared_ptr<std::vector<std::size_t>> batch_sizes(new std::vector<std::size_t>());
std::shared_ptr<std::atomic<std::size_t>> spans_received_count(new std::atomic<std::size_t>(0));
std::shared_ptr<std::atomic<bool>> is_shutdown(new std::atomic<bool>(false));
std::shared_ptr<std::atomic<std::size_t>> force_flush_counter(new std::atomic<std::size_t>(0));

auto exporter_raw = new BlockingMockSpanExporter(batch_sizes, spans_received_count, is_shutdown,
force_flush_counter);
auto exporter = std::unique_ptr<sdk::trace::SpanExporter>(exporter_raw);

sdk::trace::BatchSpanProcessorOptions options{};
options.max_export_batch_size = 100;
options.max_queue_size = 1000;

auto batch_processor = std::shared_ptr<sdk::trace::BatchSpanProcessor>(
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<std::atomic<std::size_t>> shut_down_counter(new std::atomic<std::size_t>(0));
Expand Down Expand Up @@ -220,6 +337,72 @@ TEST_F(BatchSpanProcessorTestPeer, TestForceFlush)
}
}

TEST_F(BatchSpanProcessorTestPeer, TestForceFlushDoesNotPermanentlyDrain)
{
std::shared_ptr<std::vector<std::size_t>> batch_sizes(new std::vector<std::size_t>());
std::shared_ptr<std::atomic<std::size_t>> spans_received_count(new std::atomic<std::size_t>(0));
std::shared_ptr<std::atomic<bool>> is_shutdown(new std::atomic<bool>(false));
std::shared_ptr<std::atomic<std::size_t>> force_flush_counter(new std::atomic<std::size_t>(0));
std::shared_ptr<std::atomic<std::size_t>> export_call_count(new std::atomic<std::size_t>(0));

auto exporter_raw = new BlockingMockSpanExporter(batch_sizes, spans_received_count, is_shutdown,
force_flush_counter, export_call_count);
auto exporter = std::unique_ptr<sdk::trace::SpanExporter>(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<sdk::trace::BatchSpanProcessor>(
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 (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(export_call_count->load(), 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
{
Expand Down
Loading