Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 32 additions & 2 deletions .github/scripts/graft_bench_yaml.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
"""
Expand Down Expand Up @@ -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)
Expand All @@ -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,
Expand Down Expand Up @@ -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
Expand All @@ -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"
Expand Down Expand Up @@ -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 == "*":
Expand Down
28 changes: 28 additions & 0 deletions .github/scripts/run_dcompact_bench.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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}"
Expand Down Expand Up @@ -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"
}
Expand Down Expand Up @@ -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() {
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/db_bench-dcompact-run.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
72 changes: 57 additions & 15 deletions db/compaction/compaction_job.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down Expand Up @@ -653,8 +657,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
}
Expand Down Expand Up @@ -1007,11 +1017,15 @@ try {
auto exec = exec_factory->NewExecutor(c);
std::unique_ptr<CompactionExecutor> 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;
Expand All @@ -1033,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<uint32_t>(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) {
Expand All @@ -1059,16 +1086,22 @@ 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;
}
FileDescriptor fd(file_number, path_id, min_meta.file_size,
min_meta.smallest_seqno, min_meta.largest_seqno);
Expand Down Expand Up @@ -1178,8 +1211,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();
}
Expand Down Expand Up @@ -1371,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;
}
Expand Down Expand Up @@ -2259,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();
Expand Down Expand Up @@ -2288,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
Expand Down
9 changes: 9 additions & 0 deletions db/compaction/compaction_job.h
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,8 @@ class SubcompactionState;

class CompactionJob {
public:
using FileNumberGenerator = std::function<Status(uint64_t*)>;

CompactionJob(
int job_id, Compaction* compaction, const ImmutableDBOptions& db_options,
const MutableDBOptions& mutable_db_options,
Expand Down Expand Up @@ -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().
Expand Down Expand Up @@ -315,6 +321,9 @@ class CompactionJob {
// env_option optimized for compaction table reads
FileOptions file_options_for_read_;
VersionSet* versions_;
FileNumberGenerator file_number_generator_;
std::vector<std::string> dcompact_output_files_;
bool dcompact_output_materialized_ = false;
const std::atomic<bool>* shutting_down_;
const std::atomic<bool>& manual_compaction_canceled_;
FSDirectory* db_directory_;
Expand Down
12 changes: 12 additions & 0 deletions db/db_impl/db_impl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
3 changes: 3 additions & 0 deletions db/db_impl/db_impl.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading