data_category_t category,
rewrite_gen_t gen,
SegmentProvider& sp,
- SegmentSeqAllocator &ssa)
+ SegmentSeqAllocator &ssa,
+ TokenBucket &bucket)
: store_index(store_index),
segment_allocator(nullptr, category, gen, sp, ssa),
record_submitter(crimson::common::get_conf<uint64_t>(
"seastore_journal_batch_flush_size"),
crimson::common::get_conf<double>(
"seastore_journal_batch_preferred_fullness"),
- segment_allocator)
+ segment_allocator),
+ token_bucket(bucket)
{
}
std::list<CachedExtentRef>& extents)
{
if (extents.empty()) {
- return alloc_write_iertr::now();
+ co_return;
}
- return seastar::with_gate(write_guard, [this, &t, &extents] {
- return do_write(t, extents);
+ co_await seastar::with_gate(
+ write_guard,
+ [this, &t, &extents] -> alloc_write_iertr::future<> {
+ uint64_t size = 0;
+ for (auto &e : extents) {
+ size += e->get_length();
+ }
+ try {
+ co_await trans_intr::make_interruptible(
+ token_bucket.get(size));
+ co_await do_write(t, extents);
+ } catch (...) {
+ token_bucket.release(size);
+ throw;
+ }
});
}
cold_tier_generations);
ceph_assert(dynamic_max_rewrite_generation > MIN_REWRITE_GENERATION);
+ auto main_bw_limit = crimson::common::get_conf<
+ Option::size_t>("seastore_hot_backend_bw_throttle");
+ auto secondary_bw_limit = crimson::common::get_conf<
+ Option::size_t>("seastore_cold_backend_bw_throttle");
+
+ token_buckets.emplace_back(std::make_unique<TokenBucket>(main_bw_limit));
+ token_buckets.back()->start();
+
if (trimmer->get_backend_type() == backend_type_t::SEGMENTED) {
DEBUG("initiating SegmentCleaner");
auto segment_cleaner = dynamic_cast<SegmentCleaner*>(cleaner.get());
for (rewrite_gen_t gen = OOL_GENERATION; gen < hot_tier_generations; ++gen) {
writer_refs.emplace_back(std::make_unique<SegmentedOolWriter>(store_index,
data_category_t::DATA, gen, *segment_cleaner,
- *ool_segment_seq_allocator));
+ *ool_segment_seq_allocator, *token_buckets.back()));
data_writers_by_gen[generation_to_writer(gen)] = writer_refs.back().get();
}
for (rewrite_gen_t gen = OOL_GENERATION; gen < hot_tier_generations; ++gen) {
writer_refs.emplace_back(std::make_unique<SegmentedOolWriter>(store_index,
data_category_t::METADATA, gen, *segment_cleaner,
- *ool_segment_seq_allocator));
+ *ool_segment_seq_allocator, *token_buckets.back()));
md_writers_by_gen[generation_to_writer(gen)] = writer_refs.back().get();
}
data_writers_by_gen.resize(num_writers, nullptr);
md_writers_by_gen.resize(num_writers, {});
writer_refs.emplace_back(std::make_unique<RandomBlockOolWriter>(
- rb_cleaner));
+ rb_cleaner, *token_buckets.back()));
// TODO: implement eviction in RBCleaner and introduce further writers
data_writers_by_gen[generation_to_writer(OOL_GENERATION)] = writer_refs.back().get();
md_writers_by_gen[generation_to_writer(OOL_GENERATION)] = writer_refs.back().get();
}
if (cold_cleaner) {
+ token_buckets.emplace_back(std::make_unique<TokenBucket>(secondary_bw_limit));
+ token_buckets.back()->start();
if (cold_cleaner->get_backend_type() == backend_type_t::SEGMENTED) {
auto cold_segment_cleaner = static_cast<SegmentCleaner*>(cold_cleaner.get());
for (rewrite_gen_t gen = hot_tier_generations; gen <= dynamic_max_rewrite_generation; ++gen) {
writer_refs.emplace_back(std::make_unique<SegmentedOolWriter>(store_index,
data_category_t::DATA, gen, *cold_segment_cleaner,
- *ool_segment_seq_allocator));
+ *ool_segment_seq_allocator, *token_buckets.back()));
data_writers_by_gen[generation_to_writer(gen)] = writer_refs.back().get();
}
for (rewrite_gen_t gen = hot_tier_generations; gen <= dynamic_max_rewrite_generation; ++gen) {
writer_refs.emplace_back(std::make_unique<SegmentedOolWriter>(store_index,
data_category_t::METADATA, gen, *cold_segment_cleaner,
- *ool_segment_seq_allocator));
+ *ool_segment_seq_allocator, *token_buckets.back()));
md_writers_by_gen[generation_to_writer(gen)] = writer_refs.back().get();
}
for (auto *device : cold_segment_cleaner->get_segment_manager_group()
ceph_assert(cold_cleaner->get_backend_type() == backend_type_t::RANDOM_BLOCK);
auto rb_cleaner = static_cast<RBMCleaner*>(cold_cleaner.get());
ceph_assert(rb_cleaner);
- writer_refs.emplace_back(std::make_unique<RandomBlockOolWriter>(rb_cleaner));
+ writer_refs.emplace_back(std::make_unique<RandomBlockOolWriter>(rb_cleaner, *token_buckets.back()));
for (rewrite_gen_t gen = hot_tier_generations; gen <= dynamic_max_rewrite_generation; ++gen) {
data_writers_by_gen[generation_to_writer(gen)] = writer_refs.back().get();
}
{
LOG_PREFIX(ExtentPlacementManager::close);
INFO("started");
+ for (auto &token_bucket : token_buckets) {
+ token_bucket->stop();
+ }
return crimson::do_for_each(data_writers_by_gen, [](auto &writer) {
if (writer) {
return writer->close();
std::list<CachedExtentRef>& extents)
{
if (extents.empty()) {
- return alloc_write_iertr::now();
+ co_return;
}
- return seastar::with_gate(write_guard, [this, &t, &extents] {
- seastar::lw_shared_ptr<rbm_pending_ool_t> ptr =
- seastar::make_lw_shared<rbm_pending_ool_t>();
- ptr->pending_extents = t.get_pre_alloc_list();
- assert(!t.is_conflicted());
- t.set_pending_ool(ptr);
- return do_write(t, extents
- ).finally([this, ptr=ptr] {
- if (ptr->is_conflicted) {
- for (auto &e : ptr->pending_extents) {
- rb_cleaner->mark_space_free(e->get_paddr(), e->get_length());
- }
- }
- });
+ co_await seastar::with_gate(
+ write_guard,
+ [this, &t, &extents] -> alloc_write_iertr::future<> {
+ uint64_t size = 0;
+ for (auto &extent : extents) {
+ size += extent->get_length();
+ }
+ try {
+ co_await trans_intr::make_interruptible(
+ token_bucket.get(size));
+ seastar::lw_shared_ptr<rbm_pending_ool_t> ptr =
+ seastar::make_lw_shared<rbm_pending_ool_t>();
+ ptr->pending_extents = t.get_pre_alloc_list();
+ assert(!t.is_conflicted());
+ t.set_pending_ool(ptr);
+ co_await do_write(t, extents
+ ).finally([this, ptr=ptr] {
+ if (ptr->is_conflicted) {
+ for (auto &e : ptr->pending_extents) {
+ rb_cleaner->mark_space_free(e->get_paddr(), e->get_length());
+ }
+ }
+ });
+ } catch (...) {
+ token_bucket.release(size);
+ throw;
+ }
});
}
class Cache;
+class TokenBucket {
+ struct Blocker {
+ uint64_t size = 0;
+ seastar::promise<> pr;
+ };
+public:
+ TokenBucket(uint64_t mt) :
+ tokens(mt), max_tokens(mt), timer() {}
+
+ void start() {
+ if (max_tokens != 0) {
+ tokens = max_tokens;
+ timer.set_callback([this] {
+ if (tokens == max_tokens) {
+ return;
+ }
+ assert(tokens < max_tokens);
+ tokens += std::min(
+ max_tokens / 10 + 1,
+ max_tokens - tokens);
+ do_wake();
+ });
+ if (!timer.armed()) {
+ timer.arm_periodic(std::chrono::milliseconds(100));
+ }
+ }
+ }
+
+ void stop() {
+ if (max_tokens != 0) {
+ timer.cancel();
+ tokens = std::numeric_limits<uint64_t>::max();
+ do_wake();
+ }
+ }
+
+ seastar::future<> get(uint64_t size) {
+ if (max_tokens == 0) {
+ return seastar::now();
+ }
+ if (tokens < size) {
+ size -= tokens;
+ tokens = 0;
+ blockers.emplace_back(size);
+ return blockers.back().pr.get_future();
+ } else {
+ tokens -= size;
+ return seastar::now();
+ }
+ }
+
+ void release(uint64_t size) {
+ if (tokens == max_tokens) {
+ return;
+ }
+ assert(tokens < max_tokens);
+ tokens += std::min(size, max_tokens - tokens);
+ do_wake();
+ }
+
+private:
+ void do_wake() {
+ while (!blockers.empty()) {
+ auto &next = blockers.front();
+ if (tokens < next.size) {
+ next.size -= tokens;
+ tokens = 0;
+ break;
+ } else {
+ tokens -= next.size;
+ next.pr.set_value();
+ blockers.pop_front();
+ }
+ }
+ }
+ uint64_t tokens;
+ const uint64_t max_tokens;
+ seastar::timer<seastar::steady_clock_type> timer;
+ std::list<Blocker> blockers;
+};
+
+using TokenBucketRef = std::unique_ptr<TokenBucket>;
+
/**
* ExtentOolWriter
*
data_category_t category,
rewrite_gen_t gen,
SegmentProvider &sp,
- SegmentSeqAllocator &ssa);
+ SegmentSeqAllocator &ssa,
+ TokenBucket &buckets);
backend_type_t get_type() const final {
return backend_type_t::SEGMENTED;
journal::SegmentAllocator segment_allocator;
journal::RecordSubmitter record_submitter;
seastar::gate write_guard;
+ TokenBucket &token_bucket;
};
class RandomBlockOolWriter : public ExtentOolWriter {
public:
- RandomBlockOolWriter(RBMCleaner* rb_cleaner) :
- rb_cleaner(rb_cleaner) {}
+ RandomBlockOolWriter(RBMCleaner* rb_cleaner, TokenBucket &bucket) :
+ rb_cleaner(rb_cleaner), token_bucket(bucket) {}
backend_type_t get_type() const final {
return backend_type_t::RANDOM_BLOCK;
seastar::gate write_guard;
writer_stats_t w_stats;
mutable writer_stats_t last_w_stats;
+ TokenBucket &token_bucket;
};
struct cleaner_usage_t {
std::vector<ExtentOolWriter*> data_writers_by_gen;
// gen 0 METADATA writer is the journal writer
std::vector<ExtentOolWriter*> md_writers_by_gen;
+ std::vector<TokenBucketRef> token_buckets;
std::vector<Device*> devices_by_id;
Device* primary_device = nullptr;