feat: use bulk indexing for existing messages

Signed-off-by: Zomatree <me@zomatree.live>
This commit is contained in:
Zomatree
2026-03-29 05:44:29 +01:00
parent 3fd170a7de
commit fb487130c8
3 changed files with 58 additions and 14 deletions
+28 -6
View File
@@ -1,14 +1,36 @@
use revolt_config::capture_error;
use revolt_database::{AMQP, Database};
use revolt_database::Database;
use revolt_search::ElasticsearchClient;
pub async fn index_existing_messages(db: Database, client: ElasticsearchClient) {
log::info!("Starting bulk indexing.");
pub async fn index_existing_messages(db: Database, amqp: AMQP) {
let mut generator = db.fetch_all_messages().await.expect("Database query failed");
let mut generator = db
.fetch_all_messages()
.await
.expect("Database query failed");
let mut chunk = Vec::new();
while let Some(message) = generator.next().await {
if let Err(e) = amqp.new_message_search(message).await {
log::error!("Error pushing message to RabbitMQ: {e}");
chunk.push(message);
if chunk.len() >= 1000
&& let Err(e) = client
.bulk_index_messages(&db, std::mem::take(&mut chunk))
.await
{
log::error!("Error bulk indexing messages: {e}");
capture_error(&e);
}
}
}
if !chunk.is_empty()
&& let Err(e) = client.bulk_index_messages(&db, chunk).await
{
log::error!("Error bulk indexing messages: {e}");
capture_error(&e);
}
log::info!("Finished bulk indexing.")
}