Untitled
unknown
plain_text
9 months ago
4.1 kB
18
Indexable
def _search_es_for_candidates_pairwise(self, milvus_candidates: List], keywords: List[str] = None, keyword_mode: str = "and", keyword_source: str = "all", batch_size: int = 100) -> List]:
if not self.es.is_connected():
logger.warning("Elasticsearch client not connected")
return
all_hits: List] =
normalized_candidates =
for c in milvus_candidates:
fid = c.get("frame_id")
vpath = c.get("video_path")
fp = self._make_frame_path(fid, vpath)
normalized_candidates.append({
"frame_id": fid,
"video_path": vpath,
"frame_path": fp
})
# --- BẮT ĐẦU THAY ĐỔI ---
# 1. Tạo body cho msearch
msearch_body =
frame_path_batches = # Cần lưu lại để map kết quả sau này nếu cần
for i in range(0, len(normalized_candidates), batch_size):
batch = normalized_candidates[i:i + batch_size]
frame_paths = [c["frame_path"] for c in batch if c.get("frame_path")]
if not frame_paths:
continue
frame_path_batches.append(frame_paths)
# 1a. Thêm Header
msearch_body.append({"index": self.es.index_name})
# 1b. Tạo Query Body
filter_clause = {"terms": {"frame_path": frame_paths}}
bool_query = {"bool": {"filter": [filter_clause]}}
if keywords and len(keywords) > 0:
kw_clauses =
for kw in keywords:
if keyword_source == "ocr":
kw_clauses.append({"match": {"ocr_text": {"query": kw, "fuzziness": "AUTO"}}})
elif keyword_source == "asr":
kw_clauses.append({"match": {"asr_text": {"query": kw, "fuzziness": "AUTO"}}})
else:
kw_clauses.append({
"multi_match": {
"query": kw,
"fields": ["ocr_text", "asr_text"],
"type": "best_fields",
"fuzziness": "AUTO"
}
})
if keyword_mode == "or":
bool_query["bool"]["must"] = [{"bool": {"should": kw_clauses, "minimum_should_match": 1}}]
else:
bool_query["bool"]["must"] = kw_clauses
query_body = {
"query": bool_query,
"size": len(frame_paths), # Chỉ cần lấy tối đa số lượng ID trong batch
"_source": ["frame_id", "video_path", "frame_path", "ocr_text", "asr_text", "tags", "youtube_url", "frame_idx"]
}
# 1c. Thêm Query Body
msearch_body.append(query_body)
if not msearch_body:
return
try:
# 2. Gửi MỘT request msearch DUY NHẤT
resp = self.es.es.msearch(body=msearch_body)
# 3. Lặp qua mảng 'responses'
for response_part in resp.get("responses",):
if "error" in response_part:
logger.error(f"ES msearch sub-query error: {response_part['error']}")
continue
hits = response_part.get("hits", {}).get("hits",)
for h in hits:
src = h.get("_source", {}) or {}
if src:
all_hits.append({
"frame_id": src.get("frame_id"),
"video_path": src.get("video_path"),
"frame_path": src.get("frame_path"),
"ocr_text": src.get("ocr_text", "") or "",
"asr_text": src.get("asr_text", "") or "",
"tags": src.get("tags", ""),
"youtube_url": src.get("youtube_url") or src.get("youtube"),
"frame_idx": src.get("frame_idx")
})
else:
all_hits.append(h)
except Exception as e:
logger.error(f"ES msearch query error: {e}", exc_info=True)
# --- KẾT THÚC THAY ĐỔI ---
return all_hitsEditor is loading...
Leave a Comment