perf: defer HNSW persistence during addMany() batch operations
addMany() now temporarily switches the HNSW index to deferred persist mode during the batch loop. Previously, each add() triggered ~16-20 saveHNSWData calls for modified neighbors (each a read→gzip→atomic-write cycle). For 450 items this was ~8,100 individual storage writes. With deferred mode, dirty node IDs are collected in a Set during the batch and flushed once at the end — deduplicating repeated neighbor updates. A node modified by insert #3 and again by insert #47 is written only once. Also adds setPersistMode() to TypeAwareHNSWIndex, propagating mode changes to all existing type-specific sub-indexes.
This commit is contained in:
parent
c054c41943
commit
973b6aaa39
2 changed files with 80 additions and 42 deletions
109
src/brainy.ts
109
src/brainy.ts
|
|
@ -2691,55 +2691,80 @@ export class Brainy<T = any> implements BrainyInterface<T> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Process in batches
|
// OPTIMIZATION: Defer HNSW persistence during batch insert.
|
||||||
for (let i = 0; i < params.items.length; i += batchSize) {
|
// Without this, each add() triggers ~16-20 neighbor saveHNSWData calls
|
||||||
const chunk = params.items.slice(i, i + batchSize)
|
// (each a read→gzip→atomic-write cycle). For 450 items that's ~8,100
|
||||||
|
// individual storage writes. Deferred mode collects dirty node IDs and
|
||||||
|
// flushes once at the end — deduplicating repeated neighbor updates.
|
||||||
|
const index = this.index as any
|
||||||
|
const prevPersistMode = typeof index.getPersistMode === 'function'
|
||||||
|
? index.getPersistMode()
|
||||||
|
: null
|
||||||
|
const canDefer = prevPersistMode === 'immediate'
|
||||||
|
&& typeof index.setPersistMode === 'function'
|
||||||
|
&& typeof index.flush === 'function'
|
||||||
|
|
||||||
const promises = chunk.map(async (item) => {
|
if (canDefer) {
|
||||||
try {
|
index.setPersistMode('deferred')
|
||||||
const id = await this.add(item)
|
}
|
||||||
result.successful.push(id)
|
|
||||||
} catch (error) {
|
try {
|
||||||
result.failed.push({
|
// Process in batches
|
||||||
item,
|
for (let i = 0; i < params.items.length; i += batchSize) {
|
||||||
error: (error as Error).message
|
const chunk = params.items.slice(i, i + batchSize)
|
||||||
})
|
|
||||||
if (!params.continueOnError) {
|
const promises = chunk.map(async (item) => {
|
||||||
throw error
|
try {
|
||||||
|
const id = await this.add(item)
|
||||||
|
result.successful.push(id)
|
||||||
|
} catch (error) {
|
||||||
|
result.failed.push({
|
||||||
|
item,
|
||||||
|
error: (error as Error).message
|
||||||
|
})
|
||||||
|
if (!params.continueOnError) {
|
||||||
|
throw error
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
// Parallel vs Sequential based on storage preference
|
||||||
|
if (parallel) {
|
||||||
|
await Promise.allSettled(promises)
|
||||||
|
} else {
|
||||||
|
// Sequential processing for rate-limited storage
|
||||||
|
for (const promise of promises) {
|
||||||
|
await promise
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
})
|
|
||||||
|
|
||||||
// Parallel vs Sequential based on storage preference
|
// Progress callback
|
||||||
if (parallel) {
|
if (params.onProgress) {
|
||||||
await Promise.allSettled(promises)
|
params.onProgress(
|
||||||
} else {
|
result.successful.length + result.failed.length,
|
||||||
// Sequential processing for rate-limited storage
|
result.total
|
||||||
for (const promise of promises) {
|
|
||||||
await promise
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Progress callback
|
|
||||||
if (params.onProgress) {
|
|
||||||
params.onProgress(
|
|
||||||
result.successful.length + result.failed.length,
|
|
||||||
result.total
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Adaptive delay between batches
|
|
||||||
if (i + batchSize < params.items.length && delayMs > 0) {
|
|
||||||
const batchDuration = Date.now() - lastBatchTime
|
|
||||||
|
|
||||||
// If batch was too fast, add delay to respect rate limits
|
|
||||||
if (batchDuration < delayMs) {
|
|
||||||
await new Promise(resolve =>
|
|
||||||
setTimeout(resolve, delayMs - batchDuration)
|
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
lastBatchTime = Date.now()
|
// Adaptive delay between batches
|
||||||
|
if (i + batchSize < params.items.length && delayMs > 0) {
|
||||||
|
const batchDuration = Date.now() - lastBatchTime
|
||||||
|
|
||||||
|
// If batch was too fast, add delay to respect rate limits
|
||||||
|
if (batchDuration < delayMs) {
|
||||||
|
await new Promise(resolve =>
|
||||||
|
setTimeout(resolve, delayMs - batchDuration)
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
lastBatchTime = Date.now()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
// Restore persist mode and flush all dirty nodes in one pass
|
||||||
|
if (canDefer) {
|
||||||
|
await index.flush()
|
||||||
|
index.setPersistMode(prevPersistMode)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -163,6 +163,19 @@ export class TypeAwareHNSWIndex {
|
||||||
return this.persistMode
|
return this.persistMode
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Switch persist mode at runtime. Propagates to all existing type-specific indexes.
|
||||||
|
* Used by addMany() to defer persistence during batch operations.
|
||||||
|
*/
|
||||||
|
public setPersistMode(mode: 'immediate' | 'deferred'): void {
|
||||||
|
this.persistMode = mode
|
||||||
|
for (const index of this.indexes.values()) {
|
||||||
|
if (typeof (index as any).setPersistMode === 'function') {
|
||||||
|
(index as any).setPersistMode(mode)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Get or create HNSW index for a specific type (lazy initialization)
|
* Get or create HNSW index for a specific type (lazy initialization)
|
||||||
*
|
*
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue