mirror of
https://github.com/infiniflow/ragflow.git
synced 2026-07-28 19:58:11 +08:00
Push metadata filters down to Infinity (#14974)
### What problem does this PR solve? Push metadata filters down to Infinity ### Type of change - [x] Refactoring
This commit is contained in:
@@ -1344,6 +1344,7 @@ async def search_datasets(tenant_id: str, req: dict):
|
||||
chat_mdl = LLMBundle(tenant_id, chat_model_config)
|
||||
|
||||
if meta_data_filter:
|
||||
logging.debug(f"Metadata filter: {meta_data_filter}, question: {question}, chat_mdl={'None' if chat_mdl is None else chat_mdl.llm_name}")
|
||||
local_doc_ids = await apply_meta_data_filter(
|
||||
meta_data_filter,
|
||||
None,
|
||||
|
||||
@@ -404,7 +404,7 @@ class DocMetadataService:
|
||||
)
|
||||
else:
|
||||
logging.debug(f"Backend {type(settings.docStoreConn).__name__} has no refresh_idx; skipping")
|
||||
|
||||
|
||||
logging.debug(f"Successfully inserted metadata for document {doc_id}")
|
||||
return True
|
||||
|
||||
@@ -448,7 +448,8 @@ class DocMetadataService:
|
||||
# Post-process to split combined values
|
||||
processed_meta = cls._split_combined_values(meta_fields)
|
||||
|
||||
logging.debug(f"[update_document_metadata] Updating doc_id: {doc_id}, kb_id: {kb_id}, meta_fields: {processed_meta}")
|
||||
logging.debug(
|
||||
f"[update_document_metadata] Updating doc_id: {doc_id}, kb_id: {kb_id}, meta_fields: {processed_meta}")
|
||||
|
||||
# For Elasticsearch, use efficient partial update
|
||||
if not settings.DOC_ENGINE_INFINITY and not settings.DOC_ENGINE_OCEANBASE:
|
||||
@@ -456,7 +457,8 @@ class DocMetadataService:
|
||||
index_exists = settings.docStoreConn.index_exist(index_name, "")
|
||||
if not index_exists:
|
||||
# Index doesn't exist - create it and insert directly
|
||||
logging.debug(f"[update_document_metadata] Index {index_name} does not exist, creating and inserting")
|
||||
logging.debug(
|
||||
f"[update_document_metadata] Index {index_name} does not exist, creating and inserting")
|
||||
result = settings.docStoreConn.create_doc_meta_idx(index_name)
|
||||
if result is False:
|
||||
logging.error(f"Failed to create metadata index {index_name}")
|
||||
@@ -477,7 +479,8 @@ class DocMetadataService:
|
||||
# to a backend-provided scripted assignment that fully overwrites it.
|
||||
replace_meta_fields = getattr(settings.docStoreConn, "replace_meta_fields", None)
|
||||
if callable(replace_meta_fields) and replace_meta_fields(index_name, doc_id, processed_meta):
|
||||
logging.debug(f"Successfully updated metadata for document {doc_id} via {type(settings.docStoreConn).__name__}.replace_meta_fields")
|
||||
logging.debug(
|
||||
f"Successfully updated metadata for document {doc_id} via {type(settings.docStoreConn).__name__}.replace_meta_fields")
|
||||
return True
|
||||
logging.warning(
|
||||
f"replace_meta_fields unavailable or failed on backend "
|
||||
@@ -537,7 +540,8 @@ class DocMetadataService:
|
||||
# Check if metadata table exists before attempting deletion
|
||||
# This is the key optimization - no table = no metadata = nothing to delete
|
||||
if not settings.docStoreConn.index_exist(index_name, ""):
|
||||
logging.debug(f"Metadata table {index_name} does not exist, skipping metadata deletion for document {doc_id}")
|
||||
logging.debug(
|
||||
f"Metadata table {index_name} does not exist, skipping metadata deletion for document {doc_id}")
|
||||
return True # No metadata to delete is considered success
|
||||
|
||||
# Try to get the metadata to confirm it exists before deleting
|
||||
@@ -627,7 +631,8 @@ class DocMetadataService:
|
||||
if isinstance(results, tuple) and len(results) == 2:
|
||||
# Infinity returns (DataFrame, int)
|
||||
df, total = results
|
||||
logging.debug(f"[DROP EMPTY TABLE] Infinity format - total: {total}, df length: {len(df) if hasattr(df, '__len__') else 'N/A'}")
|
||||
logging.debug(
|
||||
f"[DROP EMPTY TABLE] Infinity format - total: {total}, df length: {len(df) if hasattr(df, '__len__') else 'N/A'}")
|
||||
is_empty = (total == 0 or (hasattr(df, '__len__') and len(df) == 0))
|
||||
elif hasattr(results, 'get') and 'hits' in results:
|
||||
# ES format - MUST check this before hasattr(results, '__len__')
|
||||
@@ -791,26 +796,62 @@ class DocMetadataService:
|
||||
|
||||
@classmethod
|
||||
def filter_doc_ids_by_meta_pushdown(
|
||||
cls,
|
||||
kb_ids: List[str],
|
||||
filters: List[Dict],
|
||||
logic: str = "and",
|
||||
limit: int = 10000,
|
||||
cls,
|
||||
kb_ids: List[str],
|
||||
filters: List[Dict],
|
||||
logic: str = "and",
|
||||
limit: int = 10000,
|
||||
) -> Optional[List[str]]:
|
||||
"""Run a metadata filter directly against ES, returning matching doc IDs.
|
||||
"""Run a metadata filter directly against ES or Infinity, returning matching doc IDs.
|
||||
|
||||
Returns ``None`` to signal "push-down not viable, use the in-memory
|
||||
``meta_filter`` fallback". Reasons for ``None``:
|
||||
|
||||
- Active doc store is not Elasticsearch (Infinity / OceanBase have
|
||||
different filter semantics for the JSON ``meta_fields`` column).
|
||||
- One of the user filters cannot be expressed in ES DSL.
|
||||
- The ES request itself failed (network, mapping, missing index).
|
||||
- kb_ids or filters is empty
|
||||
- One of the user filters cannot be expressed in ES DSL or Infinity SQL
|
||||
- The request itself failed (network, mapping, missing index)
|
||||
|
||||
On success returns the deduplicated, ordered list of document IDs the
|
||||
ES query matched. Callers can union or intersect this with their own
|
||||
query matched. Callers can union or intersect this with their own
|
||||
base ``doc_ids`` rather than fetching the entire metadata table.
|
||||
"""
|
||||
if not kb_ids or not filters:
|
||||
logging.debug("Metadata filter skipped: empty kb_ids or filters")
|
||||
return None
|
||||
|
||||
try:
|
||||
kb = Knowledgebase.get_by_id(kb_ids[0])
|
||||
except Exception as e:
|
||||
logging.warning(f"Metadata filter cannot resolve tenant for kb {kb_ids[0]}: {e}")
|
||||
return None
|
||||
if not kb:
|
||||
return None
|
||||
|
||||
tenant_id = kb.tenant_id
|
||||
index_name = cls._get_doc_meta_index_name(tenant_id)
|
||||
|
||||
if not settings.docStoreConn.index_exist(index_name, ""):
|
||||
return []
|
||||
|
||||
if settings.DOC_ENGINE_INFINITY:
|
||||
return cls._filter_doc_ids_by_metadata_infinity(
|
||||
index_name, kb_ids, filters, logic
|
||||
)
|
||||
else:
|
||||
return cls._filter_doc_ids_by_metadata_es(
|
||||
index_name, kb_ids, filters, logic, limit
|
||||
)
|
||||
|
||||
@classmethod
|
||||
def _filter_doc_ids_by_metadata_es(
|
||||
cls,
|
||||
index_name: str,
|
||||
kb_ids: List[str],
|
||||
filters: List[Dict],
|
||||
logic: str,
|
||||
limit: int,
|
||||
) -> Optional[List[str]]:
|
||||
"""ES push-down path for metadata filtering."""
|
||||
from common.metadata_es_filter import (
|
||||
UnsupportedMetaFilter,
|
||||
build_meta_filter_query,
|
||||
@@ -818,14 +859,6 @@ class DocMetadataService:
|
||||
is_pushdown_supported,
|
||||
)
|
||||
|
||||
if not kb_ids:
|
||||
return []
|
||||
|
||||
if settings.DOC_ENGINE_INFINITY:
|
||||
# Infinity stores ``meta_fields`` as a JSON column without dotted
|
||||
# field access; the in-memory path is still the reliable answer.
|
||||
return None
|
||||
|
||||
es_client = getattr(settings.docStoreConn, "es", None)
|
||||
if es_client is None:
|
||||
return None
|
||||
@@ -833,35 +866,12 @@ class DocMetadataService:
|
||||
if not is_pushdown_supported(filters):
|
||||
return None
|
||||
|
||||
try:
|
||||
kb = Knowledgebase.get_by_id(kb_ids[0])
|
||||
except Exception as e:
|
||||
logging.warning(f"[meta_pushdown] cannot resolve tenant for kb {kb_ids[0]}: {e}")
|
||||
return None
|
||||
if not kb:
|
||||
return None
|
||||
|
||||
tenant_id = kb.tenant_id
|
||||
index_name = cls._get_doc_meta_index_name(tenant_id)
|
||||
|
||||
try:
|
||||
if not settings.docStoreConn.index_exist(index_name, ""):
|
||||
# No metadata index → no metadata-filtered docs. Returning an
|
||||
# empty list (rather than ``None``) so callers don't bounce
|
||||
# back to the in-memory path and re-query MySQL for nothing.
|
||||
return []
|
||||
except Exception as e:
|
||||
logging.warning(f"[meta_pushdown] index_exist check failed for {index_name}: {e}")
|
||||
return None
|
||||
|
||||
try:
|
||||
query_body = build_meta_filter_query(filters, logic, kb_ids)
|
||||
except UnsupportedMetaFilter as e:
|
||||
logging.debug(f"[meta_pushdown] falling back to in-memory: {e.reason}")
|
||||
logging.error(f"ES build query failed: {e.reason}, filters={filters}")
|
||||
return None
|
||||
|
||||
# Only the doc id is needed downstream; trimming ``_source`` keeps the
|
||||
# response small when the metadata blob is large.
|
||||
request_body = {
|
||||
**query_body,
|
||||
"size": limit,
|
||||
@@ -871,12 +881,10 @@ class DocMetadataService:
|
||||
try:
|
||||
response = es_client.search(index=index_name, body=request_body)
|
||||
except Exception as e:
|
||||
logging.warning(f"[meta_pushdown] ES query failed for {index_name}: {e}")
|
||||
logging.error(f"ES metadata filter failed for {index_name}: {e}")
|
||||
return None
|
||||
|
||||
doc_ids = extract_doc_ids(response if isinstance(response, dict) else dict(response))
|
||||
# Preserve order while removing duplicates so caller-side de-dupe stays
|
||||
# cheap.
|
||||
seen: set[str] = set()
|
||||
unique: List[str] = []
|
||||
for did in doc_ids:
|
||||
@@ -887,12 +895,52 @@ class DocMetadataService:
|
||||
|
||||
if len(unique) >= limit:
|
||||
logging.warning(
|
||||
f"[meta_pushdown] hit limit {limit} for KBs {kb_ids}; some matches may be missing"
|
||||
f"ES metadata filter hit limit {limit} for KBs {kb_ids}"
|
||||
)
|
||||
|
||||
logging.debug(f"[meta_pushdown] {len(unique)} matches for KBs {kb_ids}")
|
||||
logging.debug(f"ES metadata filter returned {len(unique)} matches for KBs {kb_ids}")
|
||||
return unique
|
||||
|
||||
@classmethod
|
||||
def _filter_doc_ids_by_metadata_infinity(
|
||||
cls,
|
||||
index_name: str,
|
||||
kb_ids: List[str],
|
||||
filters: List[Dict],
|
||||
logic: str,
|
||||
) -> Optional[List[str]]:
|
||||
"""Infinity push-down path for metadata filtering."""
|
||||
from common.metadata_infinity_filter import (
|
||||
build_infinity_filter,
|
||||
extract_doc_ids,
|
||||
is_pushdown_supported,
|
||||
)
|
||||
|
||||
if not is_pushdown_supported(filters):
|
||||
return None
|
||||
|
||||
try:
|
||||
sql_filter = build_infinity_filter(filters, logic)
|
||||
escaped_kb_ids = [k.replace("'", "''") for k in kb_ids]
|
||||
kb_filter = "kb_id IN (" + ", ".join([f"'{k}'" for k in escaped_kb_ids]) + ")"
|
||||
where_clause = f"{kb_filter} AND {sql_filter}"
|
||||
logging.debug(f"Infinity metadata filter: {where_clause}")
|
||||
|
||||
inf_conn = settings.docStoreConn.connPool.get_conn()
|
||||
try:
|
||||
db_instance = inf_conn.get_database(settings.docStoreConn.dbName)
|
||||
table_instance = db_instance.get_table(index_name)
|
||||
df, _ = table_instance.output(["id"]).filter(where_clause).to_df()
|
||||
doc_ids = extract_doc_ids(df)
|
||||
logging.debug(
|
||||
f"Infinity metadata filter returned {len(doc_ids)} doc IDs for kb_ids={kb_ids}, logic={logic}")
|
||||
return doc_ids
|
||||
finally:
|
||||
settings.docStoreConn.connPool.release_conn(inf_conn)
|
||||
except Exception:
|
||||
logging.warning("Metadata filter push-down failed; falling back to in-memory filter", exc_info=True)
|
||||
return None
|
||||
|
||||
@classmethod
|
||||
def get_metadata_keys_by_kbs(cls, kb_ids: List[str]) -> List[str]:
|
||||
"""
|
||||
@@ -955,7 +1003,8 @@ class DocMetadataService:
|
||||
if doc_meta:
|
||||
meta_mapping[doc_id] = doc_meta
|
||||
|
||||
logging.debug(f"[get_metadata_for_documents] Found metadata for {len(meta_mapping)}/{len(doc_ids) if doc_ids else 'all'} documents")
|
||||
logging.debug(
|
||||
f"[get_metadata_for_documents] Found metadata for {len(meta_mapping)}/{len(doc_ids) if doc_ids else 'all'} documents")
|
||||
return meta_mapping
|
||||
|
||||
except Exception as e:
|
||||
@@ -981,6 +1030,7 @@ class DocMetadataService:
|
||||
}
|
||||
}
|
||||
"""
|
||||
|
||||
def _is_time_string(value: str) -> bool:
|
||||
"""Check if a string value is an ISO 8601 datetime (e.g., '2026-02-03T00:00:00')."""
|
||||
if not isinstance(value, str):
|
||||
@@ -1220,7 +1270,8 @@ class DocMetadataService:
|
||||
doc_ids_set = set(doc_ids)
|
||||
missing_doc_ids = doc_ids_set - found_doc_ids
|
||||
if missing_doc_ids and updates:
|
||||
logging.debug(f"[batch_update_metadata] Inserting new metadata for documents without metadata rows: {missing_doc_ids}")
|
||||
logging.debug(
|
||||
f"[batch_update_metadata] Inserting new metadata for documents without metadata rows: {missing_doc_ids}")
|
||||
for doc_id in missing_doc_ids:
|
||||
# Apply updates to create new metadata
|
||||
meta = {}
|
||||
|
||||
Reference in New Issue
Block a user