From 01f8b89636b4f1997a7308a6c81810f9901ff79c Mon Sep 17 00:00:00 2001 From: Accusys Date: Fri, 10 Jul 2026 01:10:14 +0800 Subject: [PATCH] perf(tkg): batch INSERT for co_occurrence edges using QueryBuilder - Collect edges in Vec first, then batch insert in chunks of 100 - Uses sqlx::query_builder::QueryBuilder for efficient batch INSERT - Reduces SQL round-trips from O(edges) to O(edges/100) - Maintains ON CONFLICT DO UPDATE semantics --- src/core/processor/tkg.rs | 56 ++++++++++++++++++++------------------- 1 file changed, 29 insertions(+), 27 deletions(-) diff --git a/src/core/processor/tkg.rs b/src/core/processor/tkg.rs index af5923f..d6c2b89 100644 --- a/src/core/processor/tkg.rs +++ b/src/core/processor/tkg.rs @@ -1407,6 +1407,9 @@ async fn build_co_occurrence_edges_from_qdrant( let points = face_points; let mut edge_count = 0; + // Collect edges for batch insert + let mut edges_to_insert: Vec<(&str, i64, i64, &str, String)> = Vec::new(); + for face in points { let yolo_frame = match yolo.frames.get(&face.frame.to_string()) { Some(f) => f, @@ -1440,36 +1443,35 @@ async fn build_co_occurrence_edges_from_qdrant( "object_confidence": det.confidence, }); - if let Err(e) = sqlx::query(&format!( - r#" - INSERT INTO {} (edge_type, source_node_id, target_node_id, file_uuid, properties) - VALUES ($1, $2, $3, $4, $5::jsonb) - ON CONFLICT (file_uuid, edge_type, source_node_id, target_node_id) - DO UPDATE SET properties = COALESCE(EXCLUDED.properties, tkg_edges.properties) - "#, - edges_table - )) - .bind("CO_OCCURS_WITH") - .bind(face_node_id) - .bind(obj_node_id) - .bind(file_uuid) - .bind(serde_json::to_string(&edge_props)?) - .execute(pool) - .await - { - tracing::warn!( - "[TKG] Edge insert failed (trace={}, obj={}): {}", - face.trace_id, - det.class_name, - e - ); - continue; - } - - edge_count += 1; + edges_to_insert.push(("CO_OCCURS_WITH", face_node_id, obj_node_id, file_uuid, serde_json::to_string(&edge_props)?)); } } + // Batch insert edges + let batch_size = 100; + for chunk in edges_to_insert.chunks(batch_size) { + let mut query = sqlx::query_builder::QueryBuilder::new(&format!( + "INSERT INTO {} (edge_type, source_node_id, target_node_id, file_uuid, properties) VALUES ", + edges_table + )); + for (i, (edge_type, src, tgt, uuid, props)) in chunk.iter().enumerate() { + if i > 0 { + query.push(", "); + } + query.push("(") + .push_bind(*edge_type) + .push_bind(*src) + .push_bind(*tgt) + .push_bind(*uuid) + .push_bind(props) + .push("::jsonb)"); + } + query.push(" ON CONFLICT (file_uuid, edge_type, source_node_id, target_node_id) DO UPDATE SET properties = COALESCE(EXCLUDED.properties, tkg_edges.properties)"); + query.build().execute(pool).await?; + } + + edge_count = edges_to_insert.len(); + Ok(edge_count) }