From 1553ccc437d6bc0038808eb7bfc704eb9af24996 Mon Sep 17 00:00:00 2001 From: leipeng Date: Sun, 30 Aug 2026 15:33:16 +0800 Subject: [PATCH 1/3] deps: update rockside for HandleUpdate --- sideplugin/rockside | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sideplugin/rockside b/sideplugin/rockside index 6bdb9a4a3..8bb528ea8 160000 --- a/sideplugin/rockside +++ b/sideplugin/rockside @@ -1 +1 @@ -Subproject commit 6bdb9a4a3df4d9edc5d4d3ef12d488bb8083aa6e +Subproject commit 8bb528ea839411e00af9ddb34a13f2b5beb03e57 From 01d1c8db747aac2321e5c43ecbb0f589da048477 Mon Sep 17 00:00:00 2001 From: leipeng Date: Sun, 30 Aug 2026 16:35:07 +0800 Subject: [PATCH 2/3] fix: prevent local fallback after dcompact output materializes --- db/compaction/compaction_job.cc | 15 ++++++++++++--- db/compaction/compaction_job.h | 1 + 2 files changed, 13 insertions(+), 3 deletions(-) diff --git a/db/compaction/compaction_job.cc b/db/compaction/compaction_job.cc index 01baedd8c..87ce09e68 100644 --- a/db/compaction/compaction_job.cc +++ b/db/compaction/compaction_job.cc @@ -653,8 +653,14 @@ Status CompactionJob::Run() { } Status s = RunRemote(); if (!s.ok()) { - if (exec->AllowFallbackToLocal()) { + if (exec->AllowFallbackToLocal() && + !dcompact_output_materialized_) { s = RunLocal(); + } else if (exec->AllowFallbackToLocal()) { + ROCKS_LOG_WARN(db_options_.info_log, + "[JOB %d] Skip local fallback after dcompact output " + "materialized", + job_id_); } else { // fatal, rocksdb does not handle compact errors properly } @@ -1007,11 +1013,15 @@ try { auto exec = exec_factory->NewExecutor(c); std::unique_ptr exec_auto_del(exec); exec->SetParams(&rpc_params, c); + bool should_clean_files = false; + ROCKSDB_SCOPE_EXIT( + if (should_clean_files) { exec->CleanFiles(rpc_params, rpc_results); }); Status s = exec->Execute(rpc_params, &rpc_results); if (!s.ok()) { compact_->status = s; return s; } + should_clean_files = true; if (!rpc_results.status.ok()) { compact_->status = rpc_results.status; return rpc_results.status; @@ -1070,6 +1080,7 @@ try { compact_->status = st; return st; } + dcompact_output_materialized_ = true; FileDescriptor fd(file_number, path_id, min_meta.file_size, min_meta.smallest_seqno, min_meta.largest_seqno); FileMetaData meta; @@ -1178,8 +1189,6 @@ if (stats_) { LogFlush(db_options_.info_log); TEST_SYNC_POINT("CompactionJob::RunRemote():End"); - exec->CleanFiles(rpc_params, rpc_results); - compact_->status = Status::OK(); return Status::OK(); } diff --git a/db/compaction/compaction_job.h b/db/compaction/compaction_job.h index 1f63bfac2..a548a3bea 100644 --- a/db/compaction/compaction_job.h +++ b/db/compaction/compaction_job.h @@ -315,6 +315,7 @@ class CompactionJob { // env_option optimized for compaction table reads FileOptions file_options_for_read_; VersionSet* versions_; + bool dcompact_output_materialized_ = false; const std::atomic* shutting_down_; const std::atomic& manual_compaction_canceled_; FSDirectory* db_directory_; From 79949b8ecda10bb3d65bc1e07bece2777fa1e4df Mon Sep 17 00:00:00 2001 From: leipeng Date: Sun, 30 Aug 2026 16:50:11 +0800 Subject: [PATCH 3/3] feat: support direct dcompact output --- .github/scripts/graft_bench_yaml.py | 34 +++++++++++- .github/scripts/run_dcompact_bench.sh | 28 ++++++++++ .github/workflows/db_bench-dcompact-run.yml | 1 + db/compaction/compaction_job.cc | 59 ++++++++++++++++----- db/compaction/compaction_job.h | 8 +++ db/db_impl/db_impl.cc | 12 +++++ db/db_impl/db_impl.h | 3 ++ 7 files changed, 130 insertions(+), 15 deletions(-) diff --git a/.github/scripts/graft_bench_yaml.py b/.github/scripts/graft_bench_yaml.py index cfeabc832..6ad1b159d 100644 --- a/.github/scripts/graft_bench_yaml.py +++ b/.github/scripts/graft_bench_yaml.py @@ -5,7 +5,7 @@ machine- or per-pass fields: --set-max-background-compactions N - --worker-port / --write-buffer-size(--bytes) + --worker-port / --hoster-http-url / --write-buffer-size(--bytes) --target-file-size-base / --target-file-size-multiplier --prefix-level-writers / --fill-level-writers / --rewrite-level-writer """ @@ -114,6 +114,30 @@ def _sync_worker_port(text: str, port: int) -> str: ) +def _set_hoster_http_url(text: str, url: str) -> str: + if re.search(r"\s", url): + sys.exit("FAIL: hoster_http_url must not contain whitespace") + current = re.compile( + r"^([ \t]*hoster_http_url:[ \t]*)([^\s#]+)([ \t]*(?:#.*)?)?$", + re.MULTILINE, + ) + text, n = current.subn(rf"\g<1>{url}\g<3>", text, count=1) + if n == 1: + return text + anchor = re.compile(r"^([ \t]*)hoster_root:[^\n]*$", re.MULTILINE) + text, n = anchor.subn( + lambda match: f"{match.group(0)}\n{match.group(1)}hoster_http_url: {url}", + text, + count=1, + ) + if n != 1: + sys.exit( + "FAIL: expected exactly one hoster_root: to anchor " + f"hoster_http_url, got {n}" + ) + return text + + def main() -> None: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("yaml", type=Path) @@ -138,6 +162,7 @@ def main() -> None: help="e.g. 1 or 1.5", ) parser.add_argument("--worker-port", type=int) + parser.add_argument("--hoster-http-url") parser.add_argument( "--rewrite-level-writer", nargs=2, @@ -181,6 +206,7 @@ def main() -> None: or args.target_file_size_base is not None or args.target_file_size_multiplier is not None or args.worker_port is not None + or args.hoster_http_url is not None ) lw_ops = ( args.rewrite_level_writer is not None @@ -190,7 +216,7 @@ def main() -> None: if not scalar_ops and not lw_ops: sys.exit( "FAIL: specify at least one of --set-max-background-compactions, " - "--worker-port, --write-buffer-size[--bytes], " + "--worker-port, --hoster-http-url, --write-buffer-size[--bytes], " "--target-file-size-base, --target-file-size-multiplier, " "--prefix-level-writers, --fill-level-writers, " "--rewrite-level-writer" @@ -234,6 +260,10 @@ def main() -> None: text = _sync_worker_port(text, args.worker_port) actions.append(f"worker_port={args.worker_port}") + if args.hoster_http_url is not None: + text = _set_hoster_http_url(text, args.hoster_http_url) + actions.append(f"hoster_http_url={args.hoster_http_url}") + if args.rewrite_level_writer is not None: fro, to = args.rewrite_level_writer if fro == "*": diff --git a/.github/scripts/run_dcompact_bench.sh b/.github/scripts/run_dcompact_bench.sh index b70f319a0..26b636d48 100644 --- a/.github/scripts/run_dcompact_bench.sh +++ b/.github/scripts/run_dcompact_bench.sh @@ -18,6 +18,8 @@ # LOGDIR_BASE Parent of per-engine log dirs (default: logs) # CPU_QUOTA write-side db_bench systemd CPUQuota (default 50%) # WORKER_PORT dcompact_worker listen port (default 8080) +# HOSTER_HTTP_URL If set, allocate output file numbers from this DB host +# DcompactEtcd endpoint and write SSTs directly to cf_path # MAX_PARALLEL_COMPACTIONS (default 4) # MULTI_PROCESS ToplingZipTable: fork per compact (default 1) # ZIP_SERVER_OPTIONS ZipServer civet opts when MULTI_PROCESS=1 @@ -45,6 +47,7 @@ CPU_QUOTA="${CPU_QUOTA:-50%}" # DB path must match yaml databases.*.path. hoster_root=/dev/shm — NEVER rm -rf hoster. DB_PATH="${DB_PATH:-/dev/shm/db_bench_enterprise}" WORKER_PORT="${WORKER_PORT:-8080}" +HOSTER_HTTP_URL="${HOSTER_HTTP_URL:-}" ENGINES="${ENGINES:-zipkeyonly zipkeyvalue}" export NFS_DYNAMIC_MOUNT=0 export NFS_MOUNT_ROOT="${NFS_MOUNT_ROOT:-/dev}" @@ -155,6 +158,9 @@ make_yaml_for_engine() { if [[ -n "${WRITE_BUFFER_SIZE:-}" ]]; then graft_args+=(--write-buffer-size-bytes "$WRITE_BUFFER_SIZE") fi + if [[ -n "$HOSTER_HTTP_URL" ]]; then + graft_args+=(--hoster-http-url "$HOSTER_HTTP_URL") + fi python3 "$SCRIPT_DIR/graft_bench_yaml.py" "${graft_args[@]}" --out "$out" "$src" echo "$out" } @@ -267,6 +273,28 @@ print(int(c.get("finished",0) or 0)) fi fi echo "dcompact evidence OK (finished=${finished})" + if [[ -n "$HOSTER_HTTP_URL" ]]; then + local allocation_count + allocation_count=$(awk \ + '/DcompactEtcd allocated file number/{n++} END{print n+0}' \ + "${logdir}"/LOG-*) + if [[ "$allocation_count" -le 0 ]]; then + echo "FAIL: no DB-host file-number allocation in ${logdir}/LOG-*" >&2 + return 1 + fi + if ! grep -q '] Dcompacted ' "${logdir}"/LOG-*; then + echo "FAIL: no successfully installed direct dcompact output" >&2 + return 1 + fi + local attempt_dir + attempt_dir=$(find "$DB_PATH" -type d \ + -path "$DB_PATH/job-*/att-*" -print -quit) + if [[ -n "$attempt_dir" ]]; then + echo "FAIL: direct output left host attempt directory: $attempt_dir" >&2 + return 1 + fi + echo "direct-output evidence OK (allocations=${allocation_count}, no host attempt dirs)" + fi } prepare_db() { diff --git a/.github/workflows/db_bench-dcompact-run.yml b/.github/workflows/db_bench-dcompact-run.yml index a23fcfe90..3859bd2e4 100644 --- a/.github/workflows/db_bench-dcompact-run.yml +++ b/.github/workflows/db_bench-dcompact-run.yml @@ -142,6 +142,7 @@ jobs: export LOGDIR_BASE=logs export NUM export CPU_QUOTA + export HOSTER_HTTP_URL=http://127.0.0.1:2011/CompactionExecutorFactory/dcompact .github/scripts/run_dcompact_bench.sh - name: Run RocksDB v8.10 (write @ CPU_QUOTA + CompactionService worker) diff --git a/db/compaction/compaction_job.cc b/db/compaction/compaction_job.cc index 87ce09e68..9d0cc0690 100644 --- a/db/compaction/compaction_job.cc +++ b/db/compaction/compaction_job.cc @@ -170,6 +170,10 @@ CompactionJob::CompactionJob( file_options_for_read_( fs_->OptimizeForCompactionTableRead(file_options, db_options_)), versions_(versions), + file_number_generator_([versions](uint64_t* file_number) { + *file_number = versions->NewFileNumber(); + return Status::OK(); + }), shutting_down_(shutting_down), manual_compaction_canceled_(manual_compaction_canceled), db_directory_(db_directory), @@ -1043,6 +1047,19 @@ try { TablePropertiesCollection tp_map; auto& cf_paths = imm_cfo->cf_paths; + const bool direct_output = rpc_results.output_dir.empty(); + const uint32_t direct_output_path_id = + static_cast(cf_paths.size() - 1); + if (direct_output) { + for (const auto& files : rpc_results.output_files) { + for (const auto& file : files) { + dcompact_output_files_.push_back( + TableFileName(cf_paths, file.file_number, + direct_output_path_id)); + } + } + dcompact_output_materialized_ = !dcompact_output_files_.empty(); + } compact_->num_output_files = 0; if (rpc_results.output_files.size() != num_threads) { @@ -1069,18 +1086,23 @@ try { auto& sub_state = compact_->sub_compact_states[i]; for (const auto& min_meta : rpc_results.output_files[i]) { auto old_fnum = min_meta.file_number; - auto old_fname = MakeTableFileName(rpc_results.output_dir, old_fnum); - auto path_id = c->output_path_id(); - uint64_t file_number = versions_->NewFileNumber(); + auto path_id = direct_output ? direct_output_path_id + : c->output_path_id(); + uint64_t file_number = direct_output ? old_fnum + : versions_->NewFileNumber(); std::string new_fname = TableFileName(cf_paths, file_number, path_id); - Status st = exec->RenameFile(old_fname, new_fname, min_meta.file_size); - if (!st.ok()) { - ROCKS_LOG_ERROR(db_options_.info_log, "rename(%s, %s) = %s", - old_fname.c_str(), new_fname.c_str(), st.ToString().c_str()); - compact_->status = st; - return st; + Status st; + if (!direct_output) { + auto old_fname = MakeTableFileName(rpc_results.output_dir, old_fnum); + st = exec->RenameFile(old_fname, new_fname, min_meta.file_size); + if (!st.ok()) { + ROCKS_LOG_ERROR(db_options_.info_log, "rename(%s, %s) = %s", + old_fname.c_str(), new_fname.c_str(), st.ToString().c_str()); + compact_->status = st; + return st; + } + dcompact_output_materialized_ = true; } - dcompact_output_materialized_ = true; FileDescriptor fd(file_number, path_id, min_meta.file_size, min_meta.smallest_seqno, min_meta.largest_seqno); FileMetaData meta; @@ -1380,6 +1402,15 @@ Status CompactionJob::Install(const MutableCFOptions& mutable_cf_options, << pl_stats.bytes_written_blob; } + if (!status.ok()) { + for (const auto& file : dcompact_output_files_) { + Status delete_status = env_->DeleteFile(file); + if (!delete_status.ok() && !delete_status.IsNotFound()) { + ROCKS_LOG_WARN(db_options_.info_log, "DeleteFile(%s) = %s", + file.c_str(), delete_status.ToString().c_str()); + } + } + } CleanupCompaction(); return status; } @@ -2268,8 +2299,11 @@ Status CompactionJob::OpenCompactionOutputFile(SubcompactionState* sub_compact, CompactionOutputs& outputs) { assert(sub_compact != nullptr); - // no need to lock because VersionSet::next_file_number_ is atomic - uint64_t file_number = versions_->NewFileNumber(); + uint64_t file_number = 0; + Status s = file_number_generator_(&file_number); + if (!s.ok()) { + return s; + } std::string fname = GetTableFileName(file_number); // Fire events. ColumnFamilyData* cfd = sub_compact->compaction->column_family_data(); @@ -2297,7 +2331,6 @@ Status CompactionJob::OpenCompactionOutputFile(SubcompactionState* sub_compact, } fo_copy.temperature = temperature; - Status s; IOStatus io_s; if (IsCompactionWorker()) { // maybe s3/oss auto lazy = new LazyWritableFile(); // to avoid stat cache being stale diff --git a/db/compaction/compaction_job.h b/db/compaction/compaction_job.h index a548a3bea..ce8fd1f20 100644 --- a/db/compaction/compaction_job.h +++ b/db/compaction/compaction_job.h @@ -146,6 +146,8 @@ class SubcompactionState; class CompactionJob { public: + using FileNumberGenerator = std::function; + CompactionJob( int job_id, Compaction* compaction, const ImmutableDBOptions& db_options, const MutableDBOptions& mutable_db_options, @@ -184,6 +186,10 @@ class CompactionJob { // subcompaction results Status Run(); + void SetFileNumberGenerator(FileNumberGenerator generator) { + file_number_generator_ = std::move(generator); + } + // REQUIRED: mutex held // Add compaction input/output to the current version // Releases compaction file through Compaction::ReleaseCompactionFiles(). @@ -315,6 +321,8 @@ class CompactionJob { // env_option optimized for compaction table reads FileOptions file_options_for_read_; VersionSet* versions_; + FileNumberGenerator file_number_generator_; + std::vector dcompact_output_files_; bool dcompact_output_materialized_ = false; const std::atomic* shutting_down_; const std::atomic& manual_compaction_canceled_; diff --git a/db/db_impl/db_impl.cc b/db/db_impl/db_impl.cc index e284001c9..776518476 100644 --- a/db/db_impl/db_impl.cc +++ b/db/db_impl/db_impl.cc @@ -5764,6 +5764,18 @@ Status DBImpl::GetDbSessionId(std::string& session_id) const { return Status::OK(); } +Status DBImpl::AllocateFileNumber(const std::string& db_session_id, + uint64_t* file_number) { + if (file_number == nullptr) { + return Status::InvalidArgument("file_number is null"); + } + if (db_session_id != db_session_id_) { + return Status::InvalidArgument("db_session_id mismatch"); + } + *file_number = versions_->NewFileNumber(); + return Status::OK(); +} + namespace { SemiStructuredUniqueIdGen* DbSessionIdGen() { static SemiStructuredUniqueIdGen gen; diff --git a/db/db_impl/db_impl.h b/db/db_impl/db_impl.h index 4fb5080f7..759fa9b07 100644 --- a/db/db_impl/db_impl.h +++ b/db/db_impl/db_impl.h @@ -466,6 +466,9 @@ class DBImpl : public DB { virtual Status GetDbSessionId(std::string& session_id) const override; + Status AllocateFileNumber(const std::string& db_session_id, + uint64_t* file_number); + ColumnFamilyHandle* DefaultColumnFamily() const override final; ColumnFamilyHandle* PersistentStatsColumnFamily() const;