Skip to content

Commit 4693c1e

Browse files
zhol01825Ubuntu
authored andcommitted
Add Append RMW + Split lock DIAG histograms; merge re-queue counter fix; eval 2026-04-23 batch 1-3 results
DIAG instrumentation (IExtraSearcher.h IndexStats): - 5 atomic log2 histograms (22 buckets, 1us-1s+ / 1B-1MB+): AppendLockWait, AppendGetUs, AppendPutUs, AppendPostingBytes, SplitLockWait - HistBucketOf / HistAdd / FormatHist helpers - PrintStat unconditionally emits 5 [DIAG] lines ExtraDynamicSearcher.h: - Append single-chunk RMW path: time lock-wait, Get, Put; record bytes - Split: time write-lock acquisition - AllFinished ALL DONE branch: per-layer [DIAG] dump - MergeAsyncJob re-queue: increment m_mergeJobsInFlight + m_totalMergeSubmitted (was missing, caused in-flight underflow + completed > submitted) Eval (evaluation/2026-04-23): - benchmark_spfresh_sift1b_v10_multichunk.ini (UseMultiChunkPosting=false, 16 insert threads) - benchmark_multichunk_20260424_023309.log (live run incl. batch 1-3 DIAG output) - output.json (per-batch search/insert metrics) - ANALYSIS_BATCH1-3.md: client vs TiKV-server latency; concludes server is healthy (raw_get/put avg <1ms, 0 stalls, 78 regions); client-side overhead is the bottleneck (10-37x gap, growing per batch)
1 parent f430e35 commit 4693c1e

6 files changed

Lines changed: 12426 additions & 1 deletion

File tree

AnnService/inc/Core/SPANN/ExtraDynamicSearcher.h

Lines changed: 58 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -712,7 +712,16 @@ namespace SPTAG::SPANN {
712712
double elapsedMSeconds;
713713
{
714714
std::unique_lock<std::shared_timed_mutex> lock(m_rwLocks[headID], std::defer_lock);
715-
if (requirelock) lock.lock();
715+
if (requirelock) {
716+
// [DIAG] measure split lock wait (suspect A: lock contention)
717+
auto _lockBegin = std::chrono::high_resolution_clock::now();
718+
lock.lock();
719+
auto _lockAcq = std::chrono::high_resolution_clock::now();
720+
uint64_t _lockWaitUs = std::chrono::duration_cast<std::chrono::microseconds>(_lockAcq - _lockBegin).count();
721+
IndexStats::HistAdd(m_stat.m_splitLockWaitUs, _lockWaitUs);
722+
m_stat.m_splitLockWaitTotalUs.fetch_add(_lockWaitUs, std::memory_order_relaxed);
723+
m_stat.m_splitLockSampleCount.fetch_add(1, std::memory_order_relaxed);
724+
}
716725

717726
int retry = 0;
718727
Retry:
@@ -1233,6 +1242,8 @@ namespace SPTAG::SPANN {
12331242
if (!target.isLocal) {
12341243
if (!m_worker->SendRemoteLock(target.nodeIndex, queryResult->VID, true)) {
12351244
auto* curJob = new MergeAsyncJob(this, headID, nullptr);
1245+
m_mergeJobsInFlight++;
1246+
m_totalMergeSubmitted++;
12361247
m_splitThreadPool->add(curJob);
12371248
return ErrorCode::Success;
12381249
}
@@ -1243,13 +1254,22 @@ namespace SPTAG::SPANN {
12431254
} else {
12441255
if (!anotherLock.try_lock()) {
12451256
auto* curJob = new MergeAsyncJob(this, headID, nullptr);
1257+
m_mergeJobsInFlight++;
1258+
m_totalMergeSubmitted++;
12461259
m_splitThreadPool->add(curJob);
12471260
return ErrorCode::Success;
12481261
}
12491262
}
12501263
} else {
12511264
if (!anotherLock.try_lock()) {
12521265
auto* curJob = new MergeAsyncJob(this, headID, nullptr);
1266+
// Re-queue counts as a new submission; matched by the
1267+
// m_mergeJobsInFlight-- / m_totalMergeCompleted++ in
1268+
// MergeAsyncJob::exec(). Without these increments
1269+
// m_mergeJobsInFlight underflows to a huge uint64
1270+
// and m_totalMergeCompleted exceeds m_totalMergeSubmitted.
1271+
m_mergeJobsInFlight++;
1272+
m_totalMergeSubmitted++;
12531273
m_splitThreadPool->add(curJob);
12541274
return ErrorCode::Success;
12551275
}
@@ -1786,7 +1806,11 @@ namespace SPTAG::SPANN {
17861806
bool splitPending = false;
17871807
{
17881808
//std::shared_lock<std::shared_timed_mutex> lock(m_rwLocks[headID]); //ROCKSDB
1809+
// [DIAG] measure lock wait time (suspect A: lock contention)
1810+
auto _lockBegin = std::chrono::high_resolution_clock::now();
17891811
std::unique_lock<std::shared_timed_mutex> lock(m_rwLocks[headID]); //SPDK
1812+
auto _lockAcq = std::chrono::high_resolution_clock::now();
1813+
uint64_t _lockWaitUs = std::chrono::duration_cast<std::chrono::microseconds>(_lockAcq - _lockBegin).count();
17901814
ErrorCode ret;
17911815
if (!m_headIndex->ContainSample(headID, m_layer + 1)) {
17921816
lock.unlock();
@@ -1836,7 +1860,11 @@ namespace SPTAG::SPANN {
18361860
} else {
18371861
{ static std::atomic<int> _logOnce{0}; if (_logOnce.fetch_add(1) == 0) SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "[PATH] Append using SINGLE-KEY Get+Put path (no multi-chunk)\n"); }
18381862
std::string fullPosting;
1863+
// [DIAG] measure Get latency (suspect B/C: RMW read amplification + grpc)
1864+
auto _getBegin = std::chrono::high_resolution_clock::now();
18391865
auto getRet = db->Get(DBKey(headID), &fullPosting, MaxTimeout, &(p_exWorkSpace->m_diskRequests));
1866+
auto _getEnd = std::chrono::high_resolution_clock::now();
1867+
uint64_t _getUs = std::chrono::duration_cast<std::chrono::microseconds>(_getEnd - _getBegin).count();
18401868
if (getRet != ErrorCode::Success) fullPosting.clear();
18411869
// Diagnostic: detect stale/misaligned bytes in TiKV (e.g. residue
18421870
// from a previous run with different m_vectorInfoSize, or a prior
@@ -1851,11 +1879,25 @@ namespace SPTAG::SPANN {
18511879
}
18521880
fullPosting.append(appendPosting);
18531881
postingSize = static_cast<int>(fullPosting.size());
1882+
// [DIAG] measure Put latency + posting size
1883+
auto _putBegin = std::chrono::high_resolution_clock::now();
18541884
if ((ret = db->Put(DBKey(headID), fullPosting, MaxTimeout, &(p_exWorkSpace->m_diskRequests))) != ErrorCode::Success) {
18551885
SPTAGLIB_LOG(Helper::LogLevel::LL_Error, "Merge failed for %lld! Posting Size:%d, limit: %d\n", (std::int64_t)headID, postingSize, m_postingSizeLimit);
18561886
GetDBStats();
18571887
return ret;
18581888
}
1889+
auto _putEnd = std::chrono::high_resolution_clock::now();
1890+
uint64_t _putUs = std::chrono::duration_cast<std::chrono::microseconds>(_putEnd - _putBegin).count();
1891+
// [DIAG] record into stat histograms
1892+
IndexStats::HistAdd(m_stat.m_appendLockWaitUs, _lockWaitUs);
1893+
IndexStats::HistAdd(m_stat.m_appendGetUs, _getUs);
1894+
IndexStats::HistAdd(m_stat.m_appendPutUs, _putUs);
1895+
IndexStats::HistAdd(m_stat.m_appendPostingBytes, (uint64_t)fullPosting.size());
1896+
m_stat.m_appendLockWaitTotalUs.fetch_add(_lockWaitUs, std::memory_order_relaxed);
1897+
m_stat.m_appendGetTotalUs.fetch_add(_getUs, std::memory_order_relaxed);
1898+
m_stat.m_appendPutTotalUs.fetch_add(_putUs, std::memory_order_relaxed);
1899+
m_stat.m_appendPostingBytesTotal.fetch_add((uint64_t)fullPosting.size(), std::memory_order_relaxed);
1900+
m_stat.m_appendRmwSampleCount.fetch_add(1, std::memory_order_relaxed);
18591901
}
18601902
auto appendIOEnd = std::chrono::high_resolution_clock::now();
18611903
appendIOSeconds = std::chrono::duration_cast<std::chrono::microseconds>(appendIOEnd - appendIOBegin).count();
@@ -3038,6 +3080,21 @@ namespace SPTAG::SPANN {
30383080
m_layer, totalSplit, totalMerge, m_totalReassignSubmitted.load(), totalAppend,
30393081
m_totalSplitCompleted.load(), m_totalMergeCompleted.load(), m_totalReassignCompleted.load(),
30403082
avgSplitMs, maxSplitMs);
3083+
// [DIAG] dump diagnostic histograms (lock/RMW/grpc/byte) at every ALL DONE boundary
3084+
{
3085+
uint64_t rmwN = m_stat.m_appendRmwSampleCount.load();
3086+
uint64_t splN = m_stat.m_splitLockSampleCount.load();
3087+
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "[DIAG] layer %d %s\n", m_layer,
3088+
IndexStats::FormatHist("AppendLockWait", m_stat.m_appendLockWaitUs, m_stat.m_appendLockWaitTotalUs.load(), rmwN, "us").c_str());
3089+
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "[DIAG] layer %d %s\n", m_layer,
3090+
IndexStats::FormatHist("AppendGetUs", m_stat.m_appendGetUs, m_stat.m_appendGetTotalUs.load(), rmwN, "us").c_str());
3091+
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "[DIAG] layer %d %s\n", m_layer,
3092+
IndexStats::FormatHist("AppendPutUs", m_stat.m_appendPutUs, m_stat.m_appendPutTotalUs.load(), rmwN, "us").c_str());
3093+
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "[DIAG] layer %d %s\n", m_layer,
3094+
IndexStats::FormatHist("AppendPostBytes",m_stat.m_appendPostingBytes,m_stat.m_appendPostingBytesTotal.load(), rmwN, "B").c_str());
3095+
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "[DIAG] layer %d %s\n", m_layer,
3096+
IndexStats::FormatHist("SplitLockWait", m_stat.m_splitLockWaitUs, m_stat.m_splitLockWaitTotalUs.load(), splN, "us").c_str());
3097+
}
30413098
}
30423099
m_allDonePrinted = true;
30433100
}

AnnService/inc/Core/SPANN/IExtraSearcher.h

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,56 @@ namespace SPTAG {
111111
// GC
112112
double m_garbageCost{ 0 };
113113

114+
// ----- Diagnostic instrumentation (Append RMW + Split lock) -----
115+
// Histogram has 22 buckets covering log2 ranges:
116+
// bucket i = [2^i, 2^(i+1)) (units: microseconds for time, bytes for size)
117+
// i=0 : [1us, 2us) / [1B, 2B)
118+
// i=10 : [1ms, 2ms) / [1KB, 2KB)
119+
// i=20 : [1s, 2s) / [1MB, 2MB)
120+
// i=21 : [>=2s) / [>=2MB) (tail bucket, no upper)
121+
// Buckets are atomic so any thread can update without locking.
122+
static constexpr int kHistBuckets = 22;
123+
std::atomic_uint64_t m_appendLockWaitUs[kHistBuckets]{}; // A. lock contention
124+
std::atomic_uint64_t m_appendGetUs[kHistBuckets]{}; // B. RMW Get latency (single-chunk)
125+
std::atomic_uint64_t m_appendPutUs[kHistBuckets]{}; // B+C. RMW Put latency
126+
std::atomic_uint64_t m_appendPostingBytes[kHistBuckets]{}; // RMW posting size after Get+append
127+
std::atomic_uint64_t m_splitLockWaitUs[kHistBuckets]{}; // A. Split lock contention
128+
std::atomic_uint64_t m_appendLockWaitTotalUs{ 0 };
129+
std::atomic_uint64_t m_appendGetTotalUs{ 0 };
130+
std::atomic_uint64_t m_appendPutTotalUs{ 0 };
131+
std::atomic_uint64_t m_appendPostingBytesTotal{ 0 };
132+
std::atomic_uint64_t m_splitLockWaitTotalUs{ 0 };
133+
std::atomic_uint64_t m_appendRmwSampleCount{ 0 }; // # of RMWs that recorded above
134+
std::atomic_uint64_t m_splitLockSampleCount{ 0 };
135+
136+
static int HistBucketOf(uint64_t v) {
137+
if (v == 0) return 0;
138+
int b = 0;
139+
uint64_t x = v;
140+
while (x > 1) { x >>= 1; ++b; }
141+
if (b >= kHistBuckets) b = kHistBuckets - 1;
142+
return b;
143+
}
144+
static void HistAdd(std::atomic_uint64_t* hist, uint64_t v) {
145+
hist[HistBucketOf(v)].fetch_add(1, std::memory_order_relaxed);
146+
}
147+
static std::string FormatHist(const char* label, std::atomic_uint64_t* hist,
148+
uint64_t totalSum, uint64_t sampleCount,
149+
const char* unit) {
150+
char buf[2048];
151+
int n = snprintf(buf, sizeof(buf), " %s: count=%lu avg=%.2f%s histo[bucket=count]:",
152+
label, (unsigned long)sampleCount,
153+
sampleCount ? (double)totalSum / sampleCount : 0.0, unit);
154+
for (int i = 0; i < kHistBuckets && n < (int)sizeof(buf); ++i) {
155+
uint64_t c = hist[i].load(std::memory_order_relaxed);
156+
if (c == 0) continue;
157+
uint64_t lo = (i == 0) ? 0ULL : (1ULL << i);
158+
n += snprintf(buf + n, sizeof(buf) - n, " %lu%s+:%lu",
159+
(unsigned long)lo, (i == kHistBuckets - 1) ? ">=" : "", (unsigned long)c);
160+
}
161+
return std::string(buf);
162+
}
163+
114164
void PrintStat(int finishedInsert, bool cost = false, bool reset = false) {
115165
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "After %d insertion, head vectors split %d times, head missing %d times, same head %d times, reassign %d times, reassign scan %ld times, garbage collection %d times, merge %d times\n",
116166
finishedInsert, m_splitNum, m_headMiss.load(), m_theSameHeadNum, m_reAssignNum, m_reAssignScanNum, m_garbageNum, m_mergeNum);
@@ -131,6 +181,22 @@ namespace SPTAG {
131181
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "ReassignNum: %d, ReassignAppend TotalCost: %.3lf us, PerCost: %.3lf us\n", m_reAssignNum, m_reAssignAppendCost, m_reAssignAppendCost / m_reAssignNum);
132182
}
133183

184+
// Diagnostic histograms — independent of `cost` so they print every PrintStat
185+
{
186+
uint64_t rmwN = m_appendRmwSampleCount.load();
187+
uint64_t splN = m_splitLockSampleCount.load();
188+
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "[DIAG] %s\n",
189+
FormatHist("AppendLockWait", m_appendLockWaitUs, m_appendLockWaitTotalUs.load(), rmwN, "us").c_str());
190+
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "[DIAG] %s\n",
191+
FormatHist("AppendGetUs", m_appendGetUs, m_appendGetTotalUs.load(), rmwN, "us").c_str());
192+
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "[DIAG] %s\n",
193+
FormatHist("AppendPutUs", m_appendPutUs, m_appendPutTotalUs.load(), rmwN, "us").c_str());
194+
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "[DIAG] %s\n",
195+
FormatHist("AppendPostBytes",m_appendPostingBytes,m_appendPostingBytesTotal.load(), rmwN, "B").c_str());
196+
SPTAGLIB_LOG(Helper::LogLevel::LL_Info, "[DIAG] %s\n",
197+
FormatHist("SplitLockWait", m_splitLockWaitUs, m_splitLockWaitTotalUs.load(), splN, "us").c_str());
198+
}
199+
134200
if (reset) {
135201
m_splitNum = 0;
136202
m_headMiss = 0;

0 commit comments

Comments
 (0)