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
This commit is contained in:
+29
-27
@@ -1407,6 +1407,9 @@ async fn build_co_occurrence_edges_from_qdrant(
|
|||||||
let points = face_points;
|
let points = face_points;
|
||||||
|
|
||||||
let mut edge_count = 0;
|
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 {
|
for face in points {
|
||||||
let yolo_frame = match yolo.frames.get(&face.frame.to_string()) {
|
let yolo_frame = match yolo.frames.get(&face.frame.to_string()) {
|
||||||
Some(f) => f,
|
Some(f) => f,
|
||||||
@@ -1440,36 +1443,35 @@ async fn build_co_occurrence_edges_from_qdrant(
|
|||||||
"object_confidence": det.confidence,
|
"object_confidence": det.confidence,
|
||||||
});
|
});
|
||||||
|
|
||||||
if let Err(e) = sqlx::query(&format!(
|
edges_to_insert.push(("CO_OCCURS_WITH", face_node_id, obj_node_id, file_uuid, serde_json::to_string(&edge_props)?));
|
||||||
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;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 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)
|
Ok(edge_count)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user