batch by document and dedup loaders

This commit is contained in:
aecsocket
2026-07-10 13:26:08 +01:00
parent ca1899be92
commit 5ddae38014
5 changed files with 73 additions and 37 deletions
+1
View File
@@ -30,6 +30,7 @@ SEARCH_INCREMENTAL_INDEX_BATCH_MAX_SIZE=1000
TYPESENSE_URL=http://localhost:8108
TYPESENSE_API_KEY=modrinth
TYPESENSE_INDEX_PREFIX=labrinth
TYPESENSE_IMPORT_BATCH_SIZE=5000
REDIS_URL=redis://labrinth-redis
REDIS_MIN_CONNECTIONS=0
+1
View File
@@ -48,6 +48,7 @@ SEARCH_INCREMENTAL_INDEX_BATCH_MAX_SIZE=1000
TYPESENSE_URL=http://localhost:8108
TYPESENSE_API_KEY=modrinth
TYPESENSE_INDEX_PREFIX=labrinth
TYPESENSE_IMPORT_BATCH_SIZE=5000
REDIS_URL=redis://localhost
REDIS_MIN_CONNECTIONS=0
+1
View File
@@ -165,6 +165,7 @@ vars! {
TYPESENSE_URL: String = "http://localhost:8108";
TYPESENSE_API_KEY: String = "modrinth";
TYPESENSE_INDEX_PREFIX: String = "labrinth";
TYPESENSE_IMPORT_BATCH_SIZE: usize = 5000usize;
TYPESENSE_DELETE_BATCH_SIZE: usize = 10_000usize;
// storage
@@ -32,6 +32,7 @@ pub struct TypesenseConfig {
pub index_prefix: String,
pub meta_namespace: String,
pub index_chunk_size: i64,
pub import_batch_size: usize,
pub delete_batch_size: usize,
}
@@ -161,6 +162,7 @@ impl TypesenseConfig {
index_prefix: ENV.TYPESENSE_INDEX_PREFIX.clone(),
meta_namespace: meta_namespace.unwrap_or_default(),
index_chunk_size: ENV.SEARCH_INDEX_CHUNK_SIZE,
import_batch_size: ENV.TYPESENSE_IMPORT_BATCH_SIZE,
delete_batch_size: ENV.TYPESENSE_DELETE_BATCH_SIZE,
}
}
@@ -751,6 +753,58 @@ impl Typesense {
Ok(())
}
async fn import_document_batches(
&self,
collections: &[String],
documents: &[UploadSearchProject],
) -> Result<()> {
let batch_size = self.config.import_batch_size.max(1);
for batch in documents.chunks(batch_size) {
let jsonl = documents_to_jsonl(batch)?;
for collection in collections {
info!(
collection,
document_count = batch.len(),
content_length_bytes = jsonl.len(),
"sending Typesense document import"
);
self.client
.import_documents(collection, jsonl.clone())
.await?;
}
}
Ok(())
}
async fn existing_write_collections(&self) -> Result<Vec<String>> {
let mut collections = Vec::new();
for alias in [
self.config.get_alias_name("projects"),
self.config.get_alias_name("projects_filtered"),
] {
let live = self.client.get_alias(&alias).await?;
let shadow_alt = self.config.get_next_collection_name(&alias, true);
let shadow_current =
self.config.get_next_collection_name(&alias, false);
for collection in
live.into_iter().chain([shadow_alt, shadow_current])
{
if !collections.contains(&collection)
&& self.client.collection_exists(&collection).await?
{
collections.push(collection);
}
}
}
Ok(collections)
}
}
#[async_trait]
@@ -979,11 +1033,11 @@ impl SearchBackend for Typesense {
total += uploads.len();
cursor = next_cursor;
let jsonl = documents_to_jsonl(&uploads)?;
self.client
.import_documents(&projects_next, jsonl.clone())
.await?;
self.client.import_documents(&filtered_next, jsonl).await?;
self.import_document_batches(
&[projects_next.clone(), filtered_next.clone()],
&uploads,
)
.await?;
}
info!("swapping aliases");
@@ -1014,37 +1068,14 @@ impl SearchBackend for Typesense {
return Ok(());
}
let num_documents = documents.len();
let jsonl = documents_to_jsonl(documents)?;
for alias in [
self.config.get_alias_name("projects"),
self.config.get_alias_name("projects_filtered"),
] {
let live = self.client.get_alias(&alias).await?;
let shadow_alt = self.config.get_next_collection_name(&alias, true);
let shadow_current =
self.config.get_next_collection_name(&alias, false);
debug!(
?alias,
?live,
?shadow_alt,
?shadow_current,
num_documents,
"Inserting into alias",
);
for collection in
live.into_iter().chain([shadow_alt, shadow_current])
{
if self.client.collection_exists(&collection).await? {
debug!("Inserting into existing collection {collection:?}");
self.client
.import_documents(&collection, jsonl.clone())
.await?;
}
}
}
let collections = self.existing_write_collections().await?;
debug!(
?collections,
num_documents = documents.len(),
"Inserting into collections",
);
self.import_document_batches(&collections, documents)
.await?;
debug!("Done importing");
Ok(())
+3 -1
View File
@@ -548,10 +548,12 @@ async fn build_search_documents(
from_duplicate_version_fields(aggregated_version_fields);
// aggregated project loaders
let project_loaders = versions
let mut project_loaders = versions
.iter()
.flat_map(|x| x.loaders.clone())
.collect::<Vec<_>>();
project_loaders.sort();
project_loaders.dedup();
// all valid project types across every version of the project, so that
// filters can exclude projects that have *any* version of a given