]> git-server-git.apps.pok.os.sepia.ceph.com Git - ceph.git/commitdiff
crimson/os/seastore: throttle bandwidth to secondary devices
authorZhang Song <zhangsong02@qianxin.com>
Thu, 31 Jul 2025 08:22:28 +0000 (16:22 +0800)
committerXuehan Xu <xuxuehan@qianxin.com>
Thu, 30 Jul 2026 02:12:27 +0000 (10:12 +0800)
Signed-off-by: Zhang Song <zhangsong02@qianxin.com>
Signed-off-by: Xuehan Xu <xuxuehan@qianxin.com>
src/common/options/crimson.yaml.in
src/crimson/os/seastore/extent_placement_manager.cc
src/crimson/os/seastore/extent_placement_manager.h

index 7f4f1af71ccea0a5128e134f750c6b7b59fd2424..c4bad3917c9b1e52e706290789b99ac68b91e556 100644 (file)
@@ -255,6 +255,11 @@ options:
   level: dev
   desc: The back end used by the devices in the hot (or only) tier, valid values are SEGMENTED and RANDOM_BLOCK
   default: SEGMENTED
+- name: seastore_hot_backend_bw_throttle
+  type: size
+  level: advanced
+  desc: Size in bytes per second written to the hot devices, 0 indicates no limit
+  default: 0
 - name: seastore_cold_device_type
   type: str
   level: dev
@@ -266,6 +271,11 @@ options:
   level: dev
   desc: The backend used by cold devices (SEGMENTED or RANDOM_BLOCK)
   default: RANDOM_BLOCK
+- name: seastore_cold_backend_bw_throttle
+  type: size
+  level: advanced
+  desc: Size in bytes per second written to cold devices, 0 indicates no limit
+  default: 100_M
 - name: seastore_cbjournal_size
   type: size
   level: dev
index 2d29dbdb048e7d3a4452d186ecd68f191ba6dff7..554a479be8bf24975613d4ad7756235ed18c93c2 100644 (file)
@@ -16,7 +16,8 @@ SegmentedOolWriter::SegmentedOolWriter(
   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>(
@@ -27,7 +28,8 @@ SegmentedOolWriter::SegmentedOolWriter(
                        "seastore_journal_batch_flush_size"),
                      crimson::common::get_conf<double>(
                        "seastore_journal_batch_preferred_fullness"),
-                     segment_allocator)
+                     segment_allocator),
+    token_bucket(bucket)
 {
 }
 
@@ -182,10 +184,23 @@ SegmentedOolWriter::alloc_write_ool_extents(
   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;
+    }
   });
 }
 
@@ -207,6 +222,14 @@ void ExtentPlacementManager::init(
         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());
@@ -219,7 +242,7 @@ void ExtentPlacementManager::init(
     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();
     }
 
@@ -228,7 +251,7 @@ void ExtentPlacementManager::init(
     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();
     }
 
@@ -246,7 +269,7 @@ void ExtentPlacementManager::init(
     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();
@@ -262,18 +285,20 @@ void ExtentPlacementManager::init(
   }
 
   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()
@@ -284,7 +309,7 @@ void ExtentPlacementManager::init(
       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();
       }
@@ -594,6 +619,9 @@ ExtentPlacementManager::close()
 {
   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();
@@ -1179,22 +1207,35 @@ RandomBlockOolWriter::alloc_write_ool_extents(
   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;
+    }
   });
 }
 
index 001f2c15b46589e6db791060c2124f8cd2352530..3b94ae9938481b5b8da7ac06a310ea3a971ffb26 100644 (file)
@@ -23,6 +23,89 @@ namespace crimson::os::seastore {
 
 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
  *
@@ -75,7 +158,8 @@ public:
                      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;
@@ -129,13 +213,14 @@ private:
   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;
@@ -222,6 +307,7 @@ private:
   seastar::gate write_guard;
   writer_stats_t w_stats;
   mutable writer_stats_t last_w_stats;
+  TokenBucket &token_bucket;
 };
 
 struct cleaner_usage_t {
@@ -1215,6 +1301,7 @@ private:
   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;