From c93bdf4369f7db77f240fa9a71b979597a71ba70 Mon Sep 17 00:00:00 2001 From: Vitaliy Filippov Date: Thu, 21 May 2026 02:20:45 +0300 Subject: [PATCH] Allow degraded/misplaced writes (for better performance) --- src/osd/osd.cpp | 1 + src/osd/osd.h | 4 +- src/osd/osd_flush.cpp | 1 + src/osd/osd_peering_pg.cpp | 9 +++ src/osd/osd_peering_pg.h | 1 + src/osd/osd_primary_subops.cpp | 2 +- src/osd/osd_primary_write.cpp | 136 +++++++++++++++++++++------------ 7 files changed, 102 insertions(+), 52 deletions(-) diff --git a/src/osd/osd.cpp b/src/osd/osd.cpp index 6e5a0cf5..82448183 100644 --- a/src/osd/osd.cpp +++ b/src/osd/osd.cpp @@ -313,6 +313,7 @@ void osd_t::parse_config(bool init) pg_reshard_chunk_pause_ms = config["pg_reshard_chunk_pause_ms"].uint64_value(); if (!pg_reshard_chunk_pause_ms) pg_reshard_chunk_pause_ms = 100; + use_degraded_write = !json_is_false(config["use_degraded_write"]); if (!old_auto_scrub && auto_scrub) { // Schedule scrubbing diff --git a/src/osd/osd.h b/src/osd/osd.h index d5fd458b..69727d51 100644 --- a/src/osd/osd.h +++ b/src/osd/osd.h @@ -111,6 +111,7 @@ class osd_t bool no_rebalance = false; bool no_recovery = false; bool no_scrub = false; + bool use_degraded_write = true; bool allow_net_split = false; std::vector cfg_bind_addresses; int bind_port, listen_backlog = 128; @@ -343,6 +344,7 @@ class osd_t void continue_primary_sync(osd_op_t *cur_op); void continue_primary_del(osd_op_t *cur_op); bool check_write_queue(osd_op_t *cur_op, pg_t & pg); + bool wr_degraded(osd_op_t *cur_op); pg_osd_set_state_t* add_object_to_set(pg_t & pg, const object_id oid, const pg_osd_set_t & osd_set, uint64_t old_pg_state, int log_at_level); bool remove_object_from_state(object_id & oid, pg_osd_set_state_t **object_state, pg_t &pg, bool report = true); @@ -352,7 +354,7 @@ class osd_t osd_rmw_stripe_t *stripes, bool ref); pg_osd_set_state_t *mark_partial_write(pg_t & pg, osd_op_t *cur_op); void deref_object_state(pg_t & pg, pg_osd_set_state_t **object_state, bool deref); - bool remember_unstable_write(osd_op_t *cur_op, pg_t & pg, pg_osd_set_t & loc_set, int base_state); + bool remember_unstable_write(osd_op_t *cur_op, pg_t & pg, uint64_t *cur_set, int base_state); void handle_primary_subop(osd_op_t *subop, osd_op_t *cur_op); void handle_primary_bs_subop(osd_op_t *subop); void add_bs_subop_stats(osd_op_t *subop, bool recovery_related = false); diff --git a/src/osd/osd_flush.cpp b/src/osd/osd_flush.cpp index d537318d..de4cb285 100644 --- a/src/osd/osd_flush.cpp +++ b/src/osd/osd_flush.cpp @@ -294,6 +294,7 @@ void osd_t::submit_recovery_op(osd_recovery_op_t *op) .inode = op->oid.inode, .offset = op->oid.stripe, .len = 0, + .flags = OSD_OP_RECOVERY_RELATED, }, }; if (log_level > 2) diff --git a/src/osd/osd_peering_pg.cpp b/src/osd/osd_peering_pg.cpp index 30d79e76..4e954052 100644 --- a/src/osd/osd_peering_pg.cpp +++ b/src/osd/osd_peering_pg.cpp @@ -335,8 +335,10 @@ pg_osd_set_state_t* pg_t::add_object_to_state(const object_id oid, const uint64_ { std::vector read_target; bool found = false; + bool has_extra; uint32_t bad_mask = (LOC_OUTDATED | LOC_CORRUPTED); retry: + has_extra = false; if (scheme == POOL_SCHEME_REPLICATED) { for (auto & o: osd_set) @@ -352,6 +354,10 @@ retry: // FIXME: This is because we then use .data() and assume it's at least long read_target.resize(pg_size); } + else if (read_target.size() > pg_size) + { + has_extra = true; + } } else { @@ -360,6 +366,8 @@ retry: { if (!(o.loc_bad & bad_mask)) { + if (read_target[o.role]) + has_extra = true; read_target[o.role] = o.osd_num; found = true; } @@ -375,6 +383,7 @@ retry: state_dict[osd_set] = { .read_target = read_target, .osd_set = osd_set, + .has_extra = has_extra, .state = state, .object_count = 1, }; diff --git a/src/osd/osd_peering_pg.h b/src/osd/osd_peering_pg.h index 9c98e91a..2c60cb2d 100644 --- a/src/osd/osd_peering_pg.h +++ b/src/osd/osd_peering_pg.h @@ -29,6 +29,7 @@ struct pg_osd_set_state_t std::vector read_target; // full OSD set including additional OSDs where the object is misplaced pg_osd_set_t osd_set; + bool has_extra = false; uint64_t state = 0; uint64_t object_count = 0; uint64_t ref_count = 0; diff --git a/src/osd/osd_primary_subops.cpp b/src/osd/osd_primary_subops.cpp index 88799c5e..5fffb8f6 100644 --- a/src/osd/osd_primary_subops.cpp +++ b/src/osd/osd_primary_subops.cpp @@ -490,7 +490,7 @@ void osd_t::handle_primary_subop(osd_op_t *subop, osd_op_t *cur_op) } if ((op_data->errors + op_data->done) >= op_data->n_subops) { - if (!op_data->errors || !op_data->done || opcode != OSD_OP_SEC_WRITE && opcode != OSD_OP_SEC_WRITE_STABLE) + if (opcode != OSD_OP_SEC_WRITE && opcode != OSD_OP_SEC_WRITE_STABLE) { delete[] op_data->subops; op_data->subops = NULL; diff --git a/src/osd/osd_primary_write.cpp b/src/osd/osd_primary_write.cpp index 2139034c..5c24898b 100644 --- a/src/osd/osd_primary_write.cpp +++ b/src/osd/osd_primary_write.cpp @@ -37,6 +37,11 @@ bool osd_t::check_write_queue(osd_op_t *cur_op, pg_t & pg) return true; } +bool osd_t::wr_degraded(osd_op_t *cur_op) +{ + return use_degraded_write && !(cur_op->req.rw.flags & OSD_OP_RECOVERY_RELATED); +} + void osd_t::continue_primary_write(osd_op_t *cur_op) { if (!cur_op->op_data && !prepare_primary_rw(cur_op)) @@ -62,10 +67,13 @@ void osd_t::continue_primary_write(osd_op_t *cur_op) { return; } + if ((cur_op->req.rw.flags & OSD_OP_RECOVERY_RELATED) && cur_op->client_id != 0) + { + cur_op->req.rw.flags = 0; + } resume_1: // Determine blocks to read and write // Missing chunks are allowed to be overwritten even in incomplete objects - // FIXME: Allow to do small writes to the old (degraded/misplaced) OSD set for lower performance impact op_data->prev_set = get_object_osd_set(pg, op_data->oid, &op_data->object_state); if (op_data->object_state) { @@ -88,18 +96,23 @@ retry_1: cur_op->reply.hdr.retval = -EIO; goto continue_others; } - // Object is degraded/misplaced and will be moved to - op_data->stripes[0].read_start = 0; - op_data->stripes[0].read_end = bs_block_size; - assert(!cur_op->rmw_buf); - cur_op->rmw_buf = op_data->stripes[0].read_buf = memalign_or_die(MEM_ALIGNMENT, bs_block_size); + if (!wr_degraded(cur_op)) + { + // Object is degraded/misplaced and will be moved to + op_data->stripes[0].read_start = 0; + op_data->stripes[0].read_end = bs_block_size; + assert(!cur_op->rmw_buf); + cur_op->rmw_buf = op_data->stripes[0].read_buf = memalign_or_die(MEM_ALIGNMENT, bs_block_size); + } } } else { assert(!cur_op->rmw_buf); cur_op->rmw_buf = calc_rmw(cur_op->buf, op_data->stripes, op_data->prev_set, - pg.pg_size, pg.pg_data_size, pg.pg_cursize, pg.cur_set.data(), bs_block_size, clean_entry_bitmap_size); + pg.pg_size, pg.pg_data_size, pg.pg_cursize, + wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data(), + bs_block_size, clean_entry_bitmap_size); if (!cur_op->rmw_buf) { // Refuse partial overwrite of an incomplete object @@ -152,8 +165,7 @@ resume_3: bitmap_set(op_data->stripes[0].bmp_buf, op_data->stripes[0].write_start, op_data->stripes[0].write_end-op_data->stripes[0].write_start, bs_bitmap_granularity); // Possibly copy new data from the request into the recovery buffer - if (pg.cur_set.data() != op_data->prev_set && (op_data->stripes[0].write_start != 0 || - op_data->stripes[0].write_end != bs_block_size)) + if (cur_op->rmw_buf) { memcpy( (uint8_t*)op_data->stripes[0].read_buf + op_data->stripes[0].req_start, @@ -173,11 +185,15 @@ resume_3: // Recover missing stripes, calculate parity if (pg.scheme == POOL_SCHEME_XOR) { - calc_rmw_parity_xor(op_data->stripes, pg.pg_size, op_data->prev_set, pg.cur_set.data(), bs_block_size, clean_entry_bitmap_size); + calc_rmw_parity_xor(op_data->stripes, pg.pg_size, op_data->prev_set, + wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data(), + bs_block_size, clean_entry_bitmap_size); } else if (pg.scheme == POOL_SCHEME_EC) { - calc_rmw_parity_ec(op_data->stripes, pg.pg_size, pg.pg_data_size, op_data->prev_set, pg.cur_set.data(), bs_block_size, clean_entry_bitmap_size); + calc_rmw_parity_ec(op_data->stripes, pg.pg_size, pg.pg_data_size, op_data->prev_set, + wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data(), + bs_block_size, clean_entry_bitmap_size); } } // Send writes @@ -204,9 +220,10 @@ resume_3: { // Check that current OSD set is in history and/or add it there std::vector history_set; - for (auto peer_osd: pg.cur_set) - if (peer_osd != 0) - history_set.push_back(peer_osd); + auto new_set = wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data(); + for (int i = 0; i < pg.pg_size; i++) + if (new_set[i] != 0) + history_set.push_back(new_set[i]); std::sort(history_set.begin(), history_set.end()); auto it = std::lower_bound(pg.target_history.begin(), pg.target_history.end(), history_set); if (it == pg.target_history.end() || *it != history_set) @@ -229,7 +246,8 @@ resume_10: pg_cancel_write_queue(pg, cur_op, op_data->oid, -EPIPE); return; } - submit_primary_subops(SUBMIT_WRITE, op_data->target_ver, pg.cur_set.data(), cur_op); + submit_primary_subops(SUBMIT_WRITE, op_data->target_ver, + wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data(), cur_op); resume_4: if (op_data->n_subops > 0) { @@ -248,7 +266,8 @@ resume_5: { if (pg.scheme != POOL_SCHEME_REPLICATED) { - submit_primary_rollback_subops(cur_op, pg.cur_set.data()); + submit_primary_rollback_subops(cur_op, + wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data()); resume_11: if (op_data->n_subops > 0) { @@ -283,12 +302,24 @@ resume_12: pg_cancel_write_queue(pg, cur_op, op_data->oid, op_data->errcode); return; } + if (op_data->object_state && + op_data->object_state->has_extra && + wr_degraded(cur_op)) + { + // Mark possible extra copies which were not overwritten as outdated + mark_partial_write(pg, cur_op); + } + if (op_data->subops) + { + delete[] op_data->subops; + op_data->subops = NULL; + } if (pg.scheme != POOL_SCHEME_REPLICATED) { // Remove version override just after the write, but before stabilizing pg.ver_override.erase(op_data->oid); } - if (op_data->object_state) + if (op_data->object_state && !wr_degraded(cur_op)) { // Any kind of a non-clean object can have extra chunks, because we don't record objects // as degraded & misplaced or incomplete & misplaced at the same time. So try to remove extra chunks @@ -339,7 +370,7 @@ resume_12: } resume_6: resume_7: - if (!remember_unstable_write(cur_op, pg, pg.cur_loc_set, 6)) + if (!remember_unstable_write(cur_op, pg, wr_degraded(cur_op) ? op_data->prev_set : pg.cur_set.data(), 6)) { return; } @@ -349,7 +380,7 @@ resume_7: pg.clean_count++; pg.total_count++; } - if (op_data->object_state) + if (op_data->object_state && !wr_degraded(cur_op)) { { int recovery_type = op_data->object_state->state & (OBJ_DEGRADED|OBJ_INCOMPLETE) ? 0 : 1; @@ -467,17 +498,13 @@ void osd_t::on_change_pg_history_hook(pool_id_t pool_id, pg_num_t pg_num) } } -bool osd_t::remember_unstable_write(osd_op_t *cur_op, pg_t & pg, pg_osd_set_t & loc_set, int base_state) +bool osd_t::remember_unstable_write(osd_op_t *cur_op, pg_t & pg, uint64_t *cur_set, int base_state) { osd_primary_op_data_t *op_data = cur_op->op_data; if (op_data->st == base_state) - { goto resume_6; - } else if (op_data->st == base_state+1) - { goto resume_7; - } if (immediate_commit == IMMEDIATE_ALL) { immediate: @@ -485,24 +512,27 @@ immediate: { // Send STABILIZE ops immediately op_data->unstable_write_osds = new std::vector(); - op_data->unstable_writes = new obj_ver_id[loc_set.size()]; + op_data->unstable_writes = new obj_ver_id[pg.pg_size]; { int last_start = 0; - for (auto & chunk: loc_set) + for (int role = 0; role < pg.pg_size; role++) { - op_data->unstable_writes[last_start] = (obj_ver_id){ - .oid = { - .inode = op_data->oid.inode, - .stripe = op_data->oid.stripe | chunk.role, - }, - .version = op_data->fact_ver, - }; - op_data->unstable_write_osds->push_back((unstable_osd_num_t){ - .osd_num = chunk.osd_num, - .start = last_start, - .len = 1, - }); - last_start++; + if (cur_set[role] != 0) + { + op_data->unstable_writes[last_start] = (obj_ver_id){ + .oid = { + .inode = op_data->oid.inode, + .stripe = op_data->oid.stripe | role, + }, + .version = op_data->fact_ver, + }; + op_data->unstable_write_osds->push_back((unstable_osd_num_t){ + .osd_num = cur_set[role], + .start = last_start, + .len = 1, + }); + last_start++; + } } } submit_primary_stab_subops(cur_op); @@ -546,24 +576,30 @@ lazy: if (pg.scheme != POOL_SCHEME_REPLICATED) { // Remember version as unstable for EC/XOR - for (auto & chunk: loc_set) + for (int role = 0; role < pg.pg_size; role++) { - this->dirty_osds.insert(chunk.osd_num); - this->unstable_writes[(osd_object_id_t){ - .osd_num = chunk.osd_num, - .oid = { - .inode = op_data->oid.inode, - .stripe = op_data->oid.stripe | chunk.role, - }, - }] = op_data->fact_ver; + if (cur_set[role] != 0) + { + this->dirty_osds.insert(cur_set[role]); + this->unstable_writes[(osd_object_id_t){ + .osd_num = cur_set[role], + .oid = { + .inode = op_data->oid.inode, + .stripe = op_data->oid.stripe | role, + }, + }] = op_data->fact_ver; + } } } else { // Only remember to sync OSDs for replicated pools - for (auto & chunk: loc_set) + for (int role = 0; role < pg.pg_size; role++) { - this->dirty_osds.insert(chunk.osd_num); + if (cur_set[role] != 0) + { + this->dirty_osds.insert(cur_set[role]); + } } } // Remember PG as dirty to drop the connection when PG goes offline