Support chunked resharding in blockstores

This commit is contained in:
Vitaliy Filippov
2026-01-25 01:56:29 +03:00
parent 233d2b2a09
commit 2b801a7ffa
8 changed files with 262 additions and 77 deletions
+5 -2
View File
@@ -183,8 +183,11 @@ public:
// Update configuration
virtual void parse_config(blockstore_config_t & config) = 0;
// Reshard database for a pool
virtual void reshard(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size) = 0;
// Reshard database for a pool in chunks
// MUST be called only when nobody makes any modifications to the DB for this pool
virtual void* reshard_start(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size, uint64_t chunk_limit) = 0;
virtual bool reshard_continue(void *reshard_state, uint64_t chunk_limit) = 0;
virtual void reshard_abort(void *reshard_state) = 0;
// Event loop
virtual void loop() = 0;
+147 -36
View File
@@ -29,6 +29,15 @@
#define IMAP_MALLOC_LOW_BITS ((size_t)0x0F)
#define IMAP_MAX_LOW 16
void inode_map_put(void* & inode_idx, heap_list_item_t* li);
void inode_map_get(void *inode_idx, heap_inode_map_t::iterator & li_it, heap_list_item_t* & li, uint64_t stripe);
void inode_map_free(void* inode_idx);
bool inode_map_is_big(void* & inode_idx);
void inode_map_iterate(void* & inode_idx, std::function<void(heap_list_item_t*)> cb);
void inode_map_replace(void* & inode_idx, const heap_inode_map_t::iterator & li_it, heap_list_item_t* new_li);
void inode_map_erase(robin_hood::unordered_flat_map<inode_t, void*, i64hash_t> & pg_idx, void* & inode_idx,
const heap_inode_map_t::iterator & li_it, heap_list_item_t* li);
static inline heap_list_item_t *list_item(heap_entry_t *wr)
{
return (heap_list_item_t*)((uint8_t*)wr - offsetof(struct heap_list_item_t, entry));
@@ -658,7 +667,8 @@ void blockstore_heap_t::recheck_buffer(heap_entry_t *cwr, uint8_t *buf)
else if (!calc_checksums(cwr, buf, false))
{
// write entry is invalid, erase it and mark newer entries with garbage bit
auto & inode_idx = block_index[get_pg_id(cwr->inode, cwr->stripe)][cwr->inode];
auto & pg_idx = block_index[get_pg_id(cwr->inode, cwr->stripe)];
auto & inode_idx = pg_idx[cwr->inode];
heap_inode_map_t::iterator li_it;
heap_list_item_t *li = NULL;
inode_map_get(inode_idx, li_it, li, cwr->stripe);
@@ -684,7 +694,7 @@ void blockstore_heap_t::recheck_buffer(heap_entry_t *cwr, uint8_t *buf)
{
fprintf(stderr, "Notice: the whole object %jx:%jx only has unfinished writes, rolling back\n",
cwr->inode, cwr->stripe);
inode_map_erase(inode_idx, li_it, li);
inode_map_erase(pg_idx, inode_idx, li_it, li);
}
free_entry(li);
}
@@ -943,44 +953,138 @@ bool blockstore_heap_t::calc_block_checksums(uint32_t *block_csums, uint8_t *bit
return res;
}
void blockstore_heap_t::reshard(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size)
struct heap_reshard_state_t
{
int state = 0;
uint64_t pool_id = 0;
uint32_t old_pg_count = 0;
uint32_t pg_count = 0;
uint32_t pg_stripe_size = 0;
uint64_t chunk_size = 0;
heap_block_index_t new_shards;
heap_block_index_t old_shards;
heap_block_index_t::iterator sh_it;
robin_hood::unordered_flat_map<inode_t, void*, i64hash_t>::iterator inode_it;
heap_inode_map_t *stripe_map = NULL;
heap_inode_map_t::iterator stripe_it;
void add(heap_list_item_t *li);
bool run(uint64_t chunk_limit);
};
void heap_reshard_state_t::add(heap_list_item_t *li)
{
// like map_to_pg()
uint64_t pg_num = (li->entry.stripe / pg_stripe_size) % pg_count + 1;
uint64_t shard_id = (pool_id << (64-POOL_ID_BITS)) | pg_num;
inode_map_put(new_shards[shard_id][li->entry.inode], li);
chunk_size++;
}
bool heap_reshard_state_t::run(uint64_t chunk_limit)
{
chunk_size = 0;
if (state == 1)
goto resume_1;
else if (state == 2)
goto resume_2;
sh_it = old_shards.begin();
for (; sh_it != old_shards.end(); sh_it++)
{
inode_it = sh_it->second.begin();
for (; inode_it != sh_it->second.end(); inode_it++)
{
if (!inode_map_is_big(inode_it->second))
{
if (chunk_limit > 0 && chunk_size >= chunk_limit)
{
state = 1;
return false;
}
resume_1:
inode_map_iterate(inode_it->second, [&](heap_list_item_t *li) { add(li); });
}
else
{
stripe_map = (heap_inode_map_t*)inode_it->second;
stripe_it = stripe_map->begin();
for (; stripe_it != stripe_map->end(); stripe_it++)
{
if (chunk_limit > 0 && chunk_size >= chunk_limit)
{
state = 2;
return false;
}
resume_2:
add(*stripe_it);
}
}
inode_map_free(inode_it->second);
}
}
return true;
}
void* blockstore_heap_t::reshard_start(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size, uint64_t chunk_limit)
{
auto & pool_settings = pool_shard_settings[pool];
if (pool_settings.pg_count == pg_count && pool_settings.pg_stripe_size == pg_stripe_size)
{
return;
return NULL;
}
uint32_t old_pg_count = !pool_settings.pg_count ? 1 : pool_settings.pg_count;
uint64_t pool_id = (uint64_t)pool;
heap_block_index_t new_shards;
for (uint32_t pg_num = 0; pg_num <= old_pg_count; pg_num++)
heap_reshard_state_t *st = new heap_reshard_state_t;
st->pool_id = (uint64_t)pool;
st->pg_count = pg_count;
st->pg_stripe_size = pg_stripe_size;
st->old_pg_count = !pool_settings.pg_count ? 1 : pool_settings.pg_count;
for (uint32_t pg_num = 0; pg_num <= st->old_pg_count; pg_num++)
{
auto sh_it = block_index.find((pool_id << (64-POOL_ID_BITS)) | pg_num);
if (sh_it == block_index.end())
auto sh_it = block_index.find((st->pool_id << (64-POOL_ID_BITS)) | pg_num);
if (sh_it != block_index.end())
{
continue;
st->old_shards[pg_num] = std::move(sh_it->second);
block_index.erase(sh_it);
}
for (auto & inode_pair: sh_it->second)
{
inode_map_iterate(inode_pair.second, [&](heap_list_item_t *li)
{
// like map_to_pg()
uint64_t pg_num = (li->entry.stripe / pg_stripe_size) % pg_count + 1;
uint64_t shard_id = (pool_id << (64-POOL_ID_BITS)) | pg_num;
inode_map_put(new_shards[shard_id][li->entry.inode], li);
});
inode_map_free(inode_pair.second);
}
block_index.erase(sh_it);
}
for (auto sh_it = new_shards.begin(); sh_it != new_shards.end(); sh_it++)
bool finished = reshard_continue(st, chunk_limit);
return finished ? NULL : st;
}
bool blockstore_heap_t::reshard_continue(void *reshard_state, uint64_t chunk_limit)
{
heap_reshard_state_t *st = (heap_reshard_state_t*)reshard_state;
if (!st->run(chunk_limit))
{
return false;
}
for (auto sh_it = st->new_shards.begin(); sh_it != st->new_shards.end(); sh_it++)
{
block_index[sh_it->first] = std::move(sh_it->second);
}
pool_settings = (pool_shard_settings_t){
.pg_count = pg_count,
.pg_stripe_size = pg_stripe_size,
pool_shard_settings[st->pool_id] = (pool_shard_settings_t){
.pg_count = st->pg_count,
.pg_stripe_size = st->pg_stripe_size,
};
delete st;
return true;
}
bool blockstore_heap_t::reshard_check(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size)
{
auto set_it = pool_shard_settings.find(pool);
return (set_it != pool_shard_settings.end() &&
set_it->second.pg_count == pg_count &&
set_it->second.pg_stripe_size == pg_stripe_size);
}
void blockstore_heap_t::reshard_abort(void* reshard_state)
{
heap_reshard_state_t *st = (heap_reshard_state_t*)reshard_state;
for (auto sh_it = st->old_shards.begin(); sh_it != st->old_shards.end(); sh_it++)
{
block_index[sh_it->first] = std::move(sh_it->second);
}
delete st;
}
heap_entry_t *blockstore_heap_t::lock_and_read_entry(object_id oid)
@@ -2171,11 +2275,12 @@ void blockstore_heap_t::apply_inflight(heap_inflight_lsn_t & inflight)
if (!next)
{
assert(!prev);
auto & inode_idx = block_index[get_pg_id(wr->inode, wr->stripe)][wr->inode];
auto & pg_idx = block_index[get_pg_id(wr->inode, wr->stripe)];
auto & inode_idx = pg_idx[wr->inode];
heap_inode_map_t::iterator li_it;
heap_list_item_t *old_li = NULL;
inode_map_get(inode_idx, li_it, old_li, wr->stripe);
inode_map_erase(inode_idx, li_it, old_li);
inode_map_erase(pg_idx, inode_idx, li_it, old_li);
}
else
{
@@ -2268,7 +2373,7 @@ void blockstore_heap_t::recalc_inode_space_stats(uint64_t pool_id, bool per_inod
// This is some really crazy shit but it seems to work well :)
// At the same time it has almost zero overhead and works just as fast for fat inodes.
void blockstore_heap_t::inode_map_get(void *inode_idx, heap_inode_map_t::iterator & li_it, heap_list_item_t* & li, uint64_t stripe)
void inode_map_get(void *inode_idx, heap_inode_map_t::iterator & li_it, heap_list_item_t* & li, uint64_t stripe)
{
size_t map_n = ((size_t)inode_idx & IMAP_MALLOC_LOW_BITS);
if (!map_n)
@@ -2295,7 +2400,7 @@ void blockstore_heap_t::inode_map_get(void *inode_idx, heap_inode_map_t::iterato
}
}
void blockstore_heap_t::inode_map_free(void* inode_idx)
void inode_map_free(void* inode_idx)
{
size_t n = ((size_t)inode_idx & IMAP_MALLOC_LOW_BITS);
if (!n)
@@ -2308,7 +2413,12 @@ void blockstore_heap_t::inode_map_free(void* inode_idx)
}
}
void blockstore_heap_t::inode_map_iterate(void* & inode_idx, std::function<void(heap_list_item_t*)> cb)
bool inode_map_is_big(void* & inode_idx)
{
return !((size_t)inode_idx & IMAP_MALLOC_LOW_BITS);
}
void inode_map_iterate(void* & inode_idx, std::function<void(heap_list_item_t*)> cb)
{
size_t n = ((size_t)inode_idx & IMAP_MALLOC_LOW_BITS);
if (!n)
@@ -2335,7 +2445,7 @@ void blockstore_heap_t::inode_map_iterate(void* & inode_idx, std::function<void(
}
}
void blockstore_heap_t::inode_map_put(void* & inode_idx, heap_list_item_t* li)
void inode_map_put(void* & inode_idx, heap_list_item_t* li)
{
if (!inode_idx)
{
@@ -2404,7 +2514,7 @@ void blockstore_heap_t::inode_map_put(void* & inode_idx, heap_list_item_t* li)
}
}
void blockstore_heap_t::inode_map_replace(void* & inode_idx, const heap_inode_map_t::iterator & li_it, heap_list_item_t* new_li)
void inode_map_replace(void* & inode_idx, const heap_inode_map_t::iterator & li_it, heap_list_item_t* new_li)
{
size_t map_n = ((size_t)inode_idx & IMAP_MALLOC_LOW_BITS);
if (!map_n)
@@ -2430,7 +2540,8 @@ void blockstore_heap_t::inode_map_replace(void* & inode_idx, const heap_inode_ma
}
}
void blockstore_heap_t::inode_map_erase(void* & inode_idx, const heap_inode_map_t::iterator & li_it, heap_list_item_t* li)
void inode_map_erase(robin_hood::unordered_flat_map<inode_t, void*, i64hash_t> & pg_idx, void* & inode_idx,
const heap_inode_map_t::iterator & li_it, heap_list_item_t* li)
{
size_t map_n = ((size_t)inode_idx & IMAP_MALLOC_LOW_BITS);
if (!map_n)
@@ -2455,7 +2566,7 @@ void blockstore_heap_t::inode_map_erase(void* & inode_idx, const heap_inode_map_
else if (map_n == 1)
{
// Erase
block_index[get_pg_id(li->entry.inode, li->entry.stripe)].erase(li->entry.inode);
pg_idx.erase(li->entry.inode);
}
else
{
+7 -8
View File
@@ -137,6 +137,8 @@ struct heap_compact_t
bool do_delete;
};
struct heap_reshard_state_t;
struct heap_li_hash
{
size_t operator()(const heap_list_item_t* li) const noexcept
@@ -205,19 +207,13 @@ class blockstore_heap_t
std::function<void(bool is_data, uint64_t offset, uint64_t len, uint8_t* buf, std::function<void()>)> recheck_cb;
int recheck_queue_depth = 0;
void inode_map_put(void* & inode_idx, heap_list_item_t* li);
void inode_map_get(void *inode_idx, heap_inode_map_t::iterator & li_it, heap_list_item_t* & li, uint64_t stripe);
void inode_map_free(void* inode_idx);
void inode_map_iterate(void* & inode_idx, std::function<void(heap_list_item_t*)> cb);
void inode_map_replace(void* & inode_idx, const heap_inode_map_t::iterator & li_it, heap_list_item_t* new_li);
void inode_map_erase(void* & inode_idx, const heap_inode_map_t::iterator & li_it, heap_list_item_t* li);
uint64_t get_pg_id(inode_t inode, uint64_t stripe);
bool validate_object(heap_entry_t *obj);
void fill_recheck_queue();
int mark_used_blocks();
void recheck_buffer(heap_entry_t *cwr, uint8_t *buf);
void defragment_block(uint32_t block_num);
void reshard_add(heap_reshard_state_t *st, heap_list_item_t *li);
int allocate_entry(uint32_t entry_size, uint32_t *block_num, bool allow_last_free);
void insert_list_item(heap_list_item_t *li);
@@ -249,7 +245,10 @@ public:
// recheck small write data after reading the database from disk
bool recheck_small_writes(std::function<void(bool is_data, uint64_t offset, uint64_t len, uint8_t* buf, std::function<void()>)> read_buffer, int queue_depth);
// reshard database according to the pool's PG count
void reshard(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size);
void* reshard_start(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size, uint64_t chunk_limit);
bool reshard_continue(void* reshard_state, uint64_t chunk_limit);
bool reshard_check(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size);
void reshard_abort(void* reshard_state);
void set_no_inode_stats(const std::vector<uint64_t> & pool_ids);
void recalc_inode_space_stats(uint64_t pool_id, bool per_inode);
// read an object entry and lock it against removal
+19 -5
View File
@@ -323,9 +323,13 @@ void blockstore_impl_t::process_list(blockstore_op_t *op)
FINISH_OP(op);
return;
}
// Check if the DB needs resharding
// (we don't know about PGs from the beginning, we only create "shards" here)
heap->reshard(INODE_POOL(min_inode), pg_count, pg_stripe_size);
// Check if the DB is sharded correctly
if (!heap->reshard_check(INODE_POOL(min_inode), pg_count, pg_stripe_size))
{
op->retval = -EAGAIN;
FINISH_OP(op);
return;
}
obj_ver_id *result = NULL;
size_t stable_count = 0, unstable_count = 0;
int res = heap->list_objects(list_pg, op->min_oid, op->max_oid, &result, &stable_count, &unstable_count);
@@ -393,7 +397,17 @@ std::string blockstore_impl_t::get_op_diag(blockstore_op_t *op)
return std::string(buf);
}
void blockstore_impl_t::reshard(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size)
void* blockstore_impl_t::reshard_start(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size, uint64_t chunk_limit)
{
heap->reshard(pool, pg_count, pg_stripe_size);
return heap->reshard_start(pool, pg_count, pg_stripe_size, chunk_limit);
}
bool blockstore_impl_t::reshard_continue(void *reshard_state, uint64_t chunk_limit)
{
return heap->reshard_continue(reshard_state, chunk_limit);
}
void blockstore_impl_t::reshard_abort(void *reshard_state)
{
return heap->reshard_abort(reshard_state);
}
+3 -1
View File
@@ -189,7 +189,9 @@ public:
void parse_config(blockstore_config_t & config);
void parse_config(blockstore_config_t & config, bool init);
void reshard(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size);
void* reshard_start(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size, uint64_t chunk_limit);
bool reshard_continue(void *reshard_state, uint64_t chunk_limit);
void reshard_abort(void *reshard_state);
// Event loop
void loop();
+76 -22
View File
@@ -407,32 +407,88 @@ blockstore_clean_db_t& blockstore_impl_t::clean_db_shard(object_id oid)
return clean_db_shards[(pool_id << (64-POOL_ID_BITS)) | pg_num];
}
void blockstore_impl_t::reshard_clean_db(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size)
struct bs_reshard_state_t
{
uint64_t pool_id = (uint64_t)pool;
int state = 0;
uint64_t pool_id = 0;
uint32_t pg_count = 0;
uint32_t pg_stripe_size = 0;
uint64_t chunk_size = 0;
std::map<pool_pg_id_t, blockstore_clean_db_t> old_shards;
std::map<pool_pg_id_t, blockstore_clean_db_t> new_shards;
auto sh_it = clean_db_shards.lower_bound((pool_id << (64-POOL_ID_BITS)));
while (sh_it != clean_db_shards.end() &&
(sh_it->first >> (64-POOL_ID_BITS)) == pool_id)
std::map<pool_pg_id_t, blockstore_clean_db_t>::iterator sh_it;
blockstore_clean_db_t::iterator obj_it;
};
void* blockstore_impl_t::reshard_start(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size, uint64_t chunk_limit)
{
auto & settings = clean_db_settings[pool];
if (settings.pg_count == pg_count && settings.pg_stripe_size == pg_stripe_size)
{
for (auto & pair: sh_it->second)
{
// like map_to_pg()
uint64_t pg_num = (pair.first.stripe / pg_stripe_size) % pg_count + 1;
uint64_t shard_id = (pool_id << (64-POOL_ID_BITS)) | pg_num;
new_shards[shard_id][pair.first] = pair.second;
}
return NULL;
}
bs_reshard_state_t *st = new bs_reshard_state_t;
st->state = 0;
st->pool_id = pool;
st->pg_count = pg_count;
st->pg_stripe_size = pg_stripe_size;
auto sh_it = clean_db_shards.lower_bound((st->pool_id << (64-POOL_ID_BITS)));
while (sh_it != clean_db_shards.end() &&
(sh_it->first >> (64-POOL_ID_BITS)) == st->pool_id)
{
st->old_shards[sh_it->first] = std::move(sh_it->second);
clean_db_shards.erase(sh_it++);
}
for (sh_it = new_shards.begin(); sh_it != new_shards.end(); sh_it++)
bool finished = reshard_continue(st, chunk_limit);
return finished ? NULL : st;
}
bool blockstore_impl_t::reshard_continue(void *reshard_state, uint64_t chunk_limit)
{
bs_reshard_state_t *st = (bs_reshard_state_t*)reshard_state;
uint64_t chunk_size = 0;
if (st->state == 1)
goto resume_1;
for (st->sh_it = st->old_shards.begin(); st->sh_it != st->old_shards.end(); )
{
for (st->obj_it = st->sh_it->second.begin(); st->obj_it != st->sh_it->second.end(); st->obj_it++)
{
if (chunk_limit > 0 && chunk_size >= chunk_limit)
{
st->state = 1;
return false;
}
resume_1:
// like map_to_pg()
uint64_t pg_num = (st->obj_it->first.stripe / st->pg_stripe_size) % st->pg_count + 1;
uint64_t shard_id = (st->pool_id << (64-POOL_ID_BITS)) | pg_num;
st->new_shards[shard_id][st->obj_it->first] = st->obj_it->second;
chunk_size++;
}
st->old_shards.erase(st->sh_it++);
}
for (auto sh_it = st->new_shards.begin(); sh_it != st->new_shards.end(); sh_it++)
{
auto & to = clean_db_shards[sh_it->first];
to.swap(sh_it->second);
}
clean_db_settings[pool_id] = (pool_shard_settings_t){
.pg_count = pg_count,
.pg_stripe_size = pg_stripe_size,
clean_db_settings[st->pool_id] = (pool_shard_settings_t){
.pg_count = st->pg_count,
.pg_stripe_size = st->pg_stripe_size,
};
delete st;
return true;
}
void blockstore_impl_t::reshard_abort(void *reshard_state)
{
bs_reshard_state_t *st = (bs_reshard_state_t*)reshard_state;
for (auto sh_it = st->old_shards.begin(); sh_it != st->old_shards.end(); sh_it++)
{
auto & to = clean_db_shards[sh_it->first];
to.swap(sh_it->second);
}
delete st;
}
void blockstore_impl_t::process_list(blockstore_op_t *op)
@@ -465,7 +521,10 @@ void blockstore_impl_t::process_list(blockstore_op_t *op)
sh_it->second.pg_count != pg_count ||
sh_it->second.pg_stripe_size != pg_stripe_size)
{
reshard_clean_db(pool_id, pg_count, pg_stripe_size);
// Sharding mismatch
op->retval = -EAGAIN;
FINISH_OP(op);
return;
}
first_shard = last_shard = ((uint64_t)pool_id << (64-POOL_ID_BITS)) | list_pg;
}
@@ -807,9 +866,4 @@ std::string blockstore_impl_t::get_op_diag(blockstore_op_t *op)
return std::string(buf);
}
void blockstore_impl_t::reshard(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size)
{
reshard_clean_db(pool, pg_count, pg_stripe_size);
}
} // namespace v1
+3 -2
View File
@@ -202,7 +202,6 @@ class blockstore_impl_t: public blockstore_i
uint8_t* get_clean_entry_bitmap(uint64_t block_loc, int offset);
blockstore_clean_db_t& clean_db_shard(object_id oid);
void reshard_clean_db(pool_id_t pool_id, uint32_t pg_count, uint32_t pg_stripe_size);
void recalc_inode_space_stats(uint64_t pool_id, bool per_inode);
// Journaling
@@ -289,7 +288,9 @@ public:
void parse_config(blockstore_config_t & config, bool init);
// Reshard database for a pool
void reshard(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size);
void* reshard_start(pool_id_t pool, uint32_t pg_count, uint32_t pg_stripe_size, uint64_t chunk_limit);
bool reshard_continue(void *reshard_state, uint64_t chunk_limit);
void reshard_abort(void *reshard_state);
// Event loop
void loop();
+2 -1
View File
@@ -122,7 +122,8 @@ void osd_t::init_blockstore(std::function<void()> on_init)
// Pre-configure pool PG shards
for (auto & pool_item: st_cli.pool_config)
{
bs->reshard(pool_item.first, pool_item.second.pg_count, pool_item.second.pg_stripe_size);
auto st = bs->reshard_start(pool_item.first, pool_item.second.pg_count, pool_item.second.pg_stripe_size, 0);
assert(!st);
}
// Autosync based on the number of unstable writes to prevent stalls due to insufficient journal space
uint64_t max_autosync = bs->get_journal_size() / bs->get_block_size() / 2;