fix(knowledge): preserve Milvus schema during rebuild

This commit is contained in:
GuoQing Zhang
2026-07-14 14:59:55 +08:00
parent 1803fb38ca
commit 52d6310cd1
4 changed files with 205 additions and 128 deletions
@@ -24,14 +24,20 @@ from bisheng.sensitive_word.domain.services.sensitive_word_policy_service import
from bisheng.user.domain.models.user import UserDao
from bisheng.utils.file import download_minio_file
_WEB_LINK_SEPARATORS = ["\n\n", "\n", "", "\\.", "", ",", "", ";", "", "\\s+", ""]
_WEB_LINK_SEPARATORS = ["\n\n", "\n", "", "\\.", "", ",", "", ";", "", "\\s+", ""] # noqa: RUF001
_WEB_LINK_SEPARATOR_RULES = ["after"] * len(_WEB_LINK_SEPARATORS)
class KnowledgeFilePipeline(BaseFilePipeline):
def __init__(self, invoke_user_id: int, db_file: KnowledgeFile, preview_cache_key: str | None = None,
no_summary: bool = False, need_thumbnail: bool = False, **kwargs):
def __init__(
self,
invoke_user_id: int,
db_file: KnowledgeFile,
preview_cache_key: str | None = None,
no_summary: bool = False,
need_thumbnail: bool = False,
**kwargs,
):
split_rule = FileProcessBase(knowledge_id=db_file.knowledge_id)
if db_file.split_rule and isinstance(db_file.split_rule, str):
split_rule = FileProcessBase(**json.loads(db_file.split_rule))
@@ -83,19 +89,19 @@ class KnowledgeFilePipeline(BaseFilePipeline):
# field `abstract`". Keeping the key always present (default "") keeps
# those legacy collections writable and matches the rebuild path.
metadata.setdefault("abstract", self.db_file.abstract or "")
# Rebuilds created by older code could infer `user_metadata` as a
# required Milvus field. SQL NULL is removed by `exclude_none=True`,
# so normalize it to an empty JSON object for those collections.
metadata.setdefault("user_metadata", self.db_file.user_metadata or {})
return metadata
def prepare_local_file(self):
self.local_file_path, _ = download_minio_file(
object_name=self.db_file.object_name,
root_dir=self.tmp_dir,
calc_sha256=False
object_name=self.db_file.object_name, root_dir=self.tmp_dir, calc_sha256=False
)
def _get_image_object_dir(self) -> str | None:
return KnowledgeUtils.get_knowledge_file_image_dir(
str(self.db_file.id), self.db_file.knowledge_id
)
return KnowledgeUtils.get_knowledge_file_image_dir(str(self.db_file.id), self.db_file.knowledge_id)
def _init_abstract_transformers(self) -> list[BaseDocumentTransformer]:
if self.no_summary:
@@ -116,66 +122,84 @@ class KnowledgeFilePipeline(BaseFilePipeline):
abstract_transformers.extend(self._init_abstract_transformers())
# FileEncodingTransformer runs right after AbstractTransformer (in _init_abstract_transformers),
# using the abstract field for LLM classification. When shougang is disabled it skips internally.
abstract_transformers.append(FileEncodingTransformer(
invoke_user_id=self.invoke_user_id,
knowledge_file=self.db_file,
))
abstract_transformers.append(SimHashTransformer(knowledge_file=self.db_file))
abstract_transformers.append(ExtraFileTransformer(
loader=self.loader,
document_id=str(self.db_file.id),
knowledge_id=self.db_file.knowledge_id,
knowledge_file=self.db_file,
))
abstract_transformers.append(ImageUploadTransformer(
loader=self.loader,
document_id=str(self.db_file.id),
knowledge_id=self.db_file.knowledge_id,
retain_images=self.file_split_rule.retain_images == 1,
))
if self.need_thumbnail:
abstract_transformers.append(ThumbnailTransformer(
loader=self.loader,
abstract_transformers.append(
FileEncodingTransformer(
invoke_user_id=self.invoke_user_id,
knowledge_file=self.db_file,
))
)
)
abstract_transformers.append(SimHashTransformer(knowledge_file=self.db_file))
abstract_transformers.append(
ExtraFileTransformer(
loader=self.loader,
document_id=str(self.db_file.id),
knowledge_id=self.db_file.knowledge_id,
knowledge_file=self.db_file,
)
)
abstract_transformers.append(
ImageUploadTransformer(
loader=self.loader,
document_id=str(self.db_file.id),
knowledge_id=self.db_file.knowledge_id,
retain_images=self.file_split_rule.retain_images == 1,
)
)
if self.need_thumbnail:
abstract_transformers.append(
ThumbnailTransformer(
loader=self.loader,
knowledge_file=self.db_file,
)
)
if self.should_use_ppt_page_split():
abstract_transformers.append(DirectChunkTransformer())
elif self.should_use_hierarchical_split():
abstract_transformers.append(HierarchicalSplitterTransformer(
hierarchy_level=self.file_split_rule.hierarchy_level,
append_title=self.file_split_rule.append_title,
max_chunk_size=self.file_split_rule.max_chunk_size,
fallback_separator=self.file_split_rule.separator,
fallback_separator_rule=self.file_split_rule.separator_rule,
fallback_chunk_size=self.file_split_rule.chunk_size,
fallback_chunk_overlap=self.get_splitter_kwargs()["chunk_overlap"],
))
abstract_transformers.append(
HierarchicalSplitterTransformer(
hierarchy_level=self.file_split_rule.hierarchy_level,
append_title=self.file_split_rule.append_title,
max_chunk_size=self.file_split_rule.max_chunk_size,
fallback_separator=self.file_split_rule.separator,
fallback_separator_rule=self.file_split_rule.separator_rule,
fallback_chunk_size=self.file_split_rule.chunk_size,
fallback_chunk_overlap=self.get_splitter_kwargs()["chunk_overlap"],
)
)
else:
abstract_transformers.append(SplitterTransformer(**self.get_splitter_kwargs()))
abstract_transformers.append(PreviewCacheTransformer(
preview_cache_key=self.preview_cache_key,
file_metadata=self.file_metadata,
))
abstract_transformers.append(
PreviewCacheTransformer(
preview_cache_key=self.preview_cache_key,
file_metadata=self.file_metadata,
)
)
return abstract_transformers
def _init_excel_transformers(self) -> list[BaseDocumentTransformer]:
abstract_transformers = self._init_content_safety_transformers()
abstract_transformers.extend(self._init_abstract_transformers())
abstract_transformers.append(SimHashTransformer(knowledge_file=self.db_file))
abstract_transformers.append(ExtraFileTransformer(
loader=self.loader,
document_id=str(self.db_file.id),
knowledge_id=self.db_file.knowledge_id,
knowledge_file=self.db_file,
))
abstract_transformers.append(ImageUploadTransformer(
loader=self.loader,
document_id=str(self.db_file.id),
knowledge_id=self.db_file.knowledge_id,
retain_images=self.file_split_rule.retain_images == 1,
))
abstract_transformers.append(PreviewCacheTransformer(
preview_cache_key=self.preview_cache_key,
file_metadata=self.file_metadata,
))
abstract_transformers.append(
ExtraFileTransformer(
loader=self.loader,
document_id=str(self.db_file.id),
knowledge_id=self.db_file.knowledge_id,
knowledge_file=self.db_file,
)
)
abstract_transformers.append(
ImageUploadTransformer(
loader=self.loader,
document_id=str(self.db_file.id),
knowledge_id=self.db_file.knowledge_id,
retain_images=self.file_split_rule.retain_images == 1,
)
)
abstract_transformers.append(
PreviewCacheTransformer(
preview_cache_key=self.preview_cache_key,
file_metadata=self.file_metadata,
)
)
return abstract_transformers
@@ -1,18 +1,13 @@
from typing import List
from loguru import logger
from bisheng.common.constants.vectorstore_metadata import KNOWLEDGE_RAG_METADATA_SCHEMA
from bisheng.common.errcode import BaseErrorCode
from bisheng.common.errcode.http_error import ServerError
from bisheng.common.errcode.knowledge import KnowledgeFileFailedError
from bisheng.core.logger import trace_id_var
from bisheng.knowledge.domain.knowledge_rag import KnowledgeRag
from bisheng.knowledge.domain.models.knowledge import Knowledge, KnowledgeDao, KnowledgeState
from bisheng.knowledge.domain.models.knowledge_file import (
KnowledgeFile,
KnowledgeFileDao,
KnowledgeFileStatus
)
from bisheng.knowledge.domain.models.knowledge_file import KnowledgeFile, KnowledgeFileDao, KnowledgeFileStatus
from bisheng.knowledge.domain.services.knowledge_service import KnowledgeService
from bisheng.llm.domain import LLMService
from bisheng.worker.main import bisheng_celery
@@ -22,16 +17,16 @@ from bisheng.worker.main import bisheng_celery
def rebuild_knowledge_celery(knowledge_id: int, new_model_id: int, invoke_user_id: int) -> str:
"""
Asynchronous task to rebuild knowledge base
Args:
knowledge_id: The knowledge base uponID
new_model_id: New.. embeddingModelsID
invoke_user_id: Call UserID
Returns:
str: Task Execution Results
"""
trace_id_var.set(f'rebuild_knowledge_{knowledge_id}')
trace_id_var.set(f"rebuild_knowledge_{knowledge_id}")
logger.info(f"rebuild_knowledge_celery start knowledge_id={knowledge_id} new_model_id={new_model_id}")
try:
# Get Knowledge Base Information
@@ -42,8 +37,7 @@ def rebuild_knowledge_celery(knowledge_id: int, new_model_id: int, invoke_user_i
# 1. according knowledge_id Found knowledgefile All in the tablestatus=2Andstatus=4File, put thestatusto4
files = KnowledgeFileDao.get_files_by_multiple_status(
knowledge_id,
[KnowledgeFileStatus.SUCCESS.value, KnowledgeFileStatus.REBUILDING.value]
knowledge_id, [KnowledgeFileStatus.SUCCESS.value, KnowledgeFileStatus.REBUILDING.value]
)
# 2. According to thecollection_namewentmilvusDelete Vector Store in
KnowledgeService.delete_knowledge_file_in_vector(knowledge=knowledge, del_es=False)
@@ -76,7 +70,6 @@ def rebuild_knowledge_celery(knowledge_id: int, new_model_id: int, invoke_user_i
# 5. Update knowledge base status
if failed_files:
# DeleteesIndex andmilvusCollections to avoid data inconsistencies
_delete_es_files(knowledge, failed_files)
@@ -91,7 +84,7 @@ def rebuild_knowledge_celery(knowledge_id: int, new_model_id: int, invoke_user_i
return f"knowledge {knowledge_id} rebuild completed"
except Exception as e:
logger.exception(f"rebuild_knowledge_celery error: {str(e)}")
logger.exception(f"rebuild_knowledge_celery error: {e!s}")
# Unexpected handles during asynchronous tasksknowledgeSet to4
try:
knowledge = KnowledgeDao.query_by_id(knowledge_id)
@@ -99,12 +92,12 @@ def rebuild_knowledge_celery(knowledge_id: int, new_model_id: int, invoke_user_i
knowledge.state = KnowledgeState.FAILED.value
KnowledgeDao.update_one(knowledge)
except Exception as e2:
logger.exception(f"Failed to update knowledge state after error: {str(e2)}")
logger.exception(f"Failed to update knowledge state after error: {e2!s}")
raise e
def _delete_es_files(knowledge: Knowledge, file_ids: List[int]):
def _delete_es_files(knowledge: Knowledge, file_ids: list[int]):
"""DeleteESFile data in"""
try:
index_name = knowledge.index_name or knowledge.collection_name
@@ -115,23 +108,18 @@ def _delete_es_files(knowledge: Knowledge, file_ids: List[int]):
return
for file_id in file_ids:
delete_query = {
"query": {
"match": {
"metadata.document_id": file_id
}
}
}
delete_query = {"query": {"match": {"metadata.document_id": file_id}}}
response = es_client.client.delete_by_query(index=index_name, body=delete_query)
deleted = response.get("deleted", 0)
logger.info(f"Deleted {deleted} documents from ES for file_id={file_id}")
except Exception as e:
logger.exception(f"Failed to delete ES files for knowledge_id={knowledge.id}: {str(e)}")
logger.exception(f"Failed to delete ES files for knowledge_id={knowledge.id}: {e!s}")
def _rebuild_embeddings(knowledge: Knowledge, files: List[KnowledgeFile], new_model_id: int, invoke_user_id: int) -> \
tuple[List[int], List[int]]:
def _rebuild_embeddings(
knowledge: Knowledge, files: list[KnowledgeFile], new_model_id: int, invoke_user_id: int
) -> tuple[list[int], list[int]]:
"""
Rebuildembeddings
@@ -149,24 +137,30 @@ def _rebuild_embeddings(knowledge: Knowledge, files: List[KnowledgeFile], new_mo
# Get newembeddingModel and createMilvusClient
logger.info(f"[DEBUG] Begin initializing newembeddingModelsmodel_id={new_model_id}")
new_embeddings = LLMService.get_bisheng_knowledge_embedding_sync(model_id=new_model_id,
invoke_user_id=invoke_user_id)
new_embeddings = LLMService.get_bisheng_knowledge_embedding_sync(
model_id=new_model_id, invoke_user_id=invoke_user_id
)
logger.info(
f"[DEBUG] Slider Created Successfully.embeddingModel Instance: {type(new_embeddings).__name__}, model_id={getattr(new_embeddings, 'model_id', 'unknown')}")
f"[DEBUG] Slider Created Successfully.embeddingModel Instance: {type(new_embeddings).__name__}, model_id={getattr(new_embeddings, 'model_id', 'unknown')}"
)
# TestembeddingIs the model available
try:
test_result = new_embeddings.embed_query("Test text")
logger.info(
f"[DEBUG] EmbeddingModel tested successfully, dimension returned: {len(test_result) if test_result else 'None'}")
f"[DEBUG] EmbeddingModel tested successfully, dimension returned: {len(test_result) if test_result else 'None'}"
)
except Exception as e:
logger.error(f"[DEBUG] EmbeddingModel Test Failed: {str(e)}")
logger.error(f"[DEBUG] EmbeddingModel Test Failed: {e!s}")
# Model test failure should terminate the entire process, not continue
raise Exception(f"EmbeddingModel not available: {str(e)}")
raise Exception(f"EmbeddingModel not available: {e!s}")
vector_client = KnowledgeRag.init_knowledge_milvus_vectorstore_sync(invoke_user_id=invoke_user_id,
knowledge=knowledge,
embeddings=new_embeddings)
vector_client = KnowledgeRag.init_knowledge_milvus_vectorstore_sync(
invoke_user_id=invoke_user_id,
knowledge=knowledge,
embeddings=new_embeddings,
metadata_schemas=KNOWLEDGE_RAG_METADATA_SCHEMA,
)
logger.info(f"[DEBUG] Slider Created Successfully.MilvusClientcollection_name={knowledge.collection_name}")
# OthersESWhether the index is present (check in advance, avoid double-checking in the loop)
@@ -186,11 +180,11 @@ def _rebuild_embeddings(knowledge: Knowledge, files: List[KnowledgeFile], new_mo
else:
failed_files.append(file.id)
except Exception as e:
logger.exception(f"Failed to rebuild embeddings for file_id={file.id}: {str(e)}")
logger.exception(f"Failed to rebuild embeddings for file_id={file.id}: {e!s}")
failed_files.append(file.id)
except Exception as e:
logger.exception(f"Failed to rebuild embeddings: {str(e)}")
logger.exception(f"Failed to rebuild embeddings: {e!s}")
# If the entire process fails, all unsuccessful files are marked as failed
failed_files.extend([f.id for f in files if f.id not in success_files])
@@ -202,14 +196,7 @@ def _process_single_file(file, es_client, index_name, vector_client):
logger.info(f"Rebuilding embeddings for file_id={file.id}")
# FROMESGet all of this file inchunks
search_query = {
"query": {
"match": {
"metadata.document_id": file.id
}
},
"size": 10000
}
search_query = {"query": {"match": {"metadata.document_id": file.id}}, "size": 10000}
logger.debug(f"ES search query: {search_query}")
@@ -237,23 +224,20 @@ def _process_single_file(file, es_client, index_name, vector_client):
logger.info(f"Found {len(texts)} chunks for file_id={file.id}")
# Insert data intoMilvus
logger.info(f"[DEBUG] Upcoming Callsvector_client.add_textstextsQuantity={len(texts)}")
logger.info(f"[DEBUG] Upcoming Callsvector_client.add_texts, textsQuantity={len(texts)}")
logger.info(f"[DEBUG] First text example: {texts[0][:100] if texts else 'No texts'}...")
try:
vector_client.add_texts(texts=texts, metadatas=metadatas)
logger.info(f"[DEBUG] vector_client.add_textsCall successful")
logger.info("[DEBUG] vector_client.add_textsCall successful")
return True
except Exception as add_error:
logger.error(f"[DEBUG] vector_client.add_textsCall failed: {str(add_error)}")
logger.error(f"[DEBUG] vector_client.add_textsCall failed: {add_error!s}")
raise add_error
def get_all_es_chunks(es_client, index_name, query):
result = es_client.search(index=index_name,
body=query,
size=5000,
scroll="1m")
result = es_client.search(index=index_name, body=query, size=5000, scroll="1m")
res = []
def handle_hits(hits):
@@ -261,14 +245,14 @@ def get_all_es_chunks(es_client, index_name, query):
res.append(hit)
handle_hits(result.get("hits", {}).get("hits", []))
scroll_id = result.get('_scroll_id')
scroll_id = result.get("_scroll_id")
while scroll_id:
result = es_client.scroll(scroll_id=scroll_id, scroll='1m')
tmp_hits = result.get('hits', {}).get('hits', [])
result = es_client.scroll(scroll_id=scroll_id, scroll="1m")
tmp_hits = result.get("hits", {}).get("hits", [])
if not tmp_hits:
break
handle_hits(tmp_hits)
scroll_id = result.get('_scroll_id')
scroll_id = result.get("_scroll_id")
if scroll_id:
es_client.clear_scroll(scroll_id=scroll_id)
return res
@@ -287,9 +271,10 @@ def rebuild_knowledge_file_chunk(file_id: int):
except BaseErrorCode as e:
KnowledgeFileDao.update_file_status([db_file.id], KnowledgeFileStatus.FAILED, e.to_json_str())
except Exception as e:
logger.exception(f"Failed to rebuild knowledge file chunk: {str(e)}")
KnowledgeFileDao.update_file_status([db_file.id], KnowledgeFileStatus.FAILED,
ServerError(exception=e).to_json_str())
logger.exception(f"Failed to rebuild knowledge file chunk: {e!s}")
KnowledgeFileDao.update_file_status(
[db_file.id], KnowledgeFileStatus.FAILED, ServerError(exception=e).to_json_str()
)
def _rebuild_knowledge_file_chunk(db_file: KnowledgeFile):
@@ -298,17 +283,11 @@ def _rebuild_knowledge_file_chunk(db_file: KnowledgeFile):
es_client = KnowledgeRag.init_knowledge_es_vectorstore_sync(db_knowledge)
index_name = db_knowledge.index_name or db_knowledge.collection_name
query = {
"query": {
"match": {
"metadata.document_id": db_file.id
}
}
}
query = {"query": {"match": {"metadata.document_id": db_file.id}}}
chunks = get_all_es_chunks(es_client.client, index_name, query)
if not chunks:
logger.warning(f"No chunks found for")
logger.warning("No chunks found for")
return
logger.info(f"Found {len(chunks)} chunks in ES for file_id={db_file.id}")
@@ -351,7 +330,7 @@ def _rebuild_knowledge_file_chunk(db_file: KnowledgeFile):
try:
milvus_client.col.delete(f"pk in {pks_to_delete}")
except Exception as e:
logger.warning(f"Failed to delete old pk(s) from Milvus: {str(e)}")
logger.warning(f"Failed to delete old pk(s) from Milvus: {e!s}")
# Re-insert into Milvus and ES
logger.info(f"Re-inserting {len(texts)} chunks for file_id={db_file.id} into vector stores")
@@ -0,0 +1,29 @@
from datetime import datetime
from types import SimpleNamespace
from bisheng.knowledge.rag.knowledge_file_pipeline import KnowledgeFilePipeline, UserDao
def test_file_metadata_normalizes_null_user_metadata(monkeypatch):
monkeypatch.setattr(
UserDao,
"get_user",
lambda _user_id: SimpleNamespace(user_name="tester"),
)
pipeline = object.__new__(KnowledgeFilePipeline)
pipeline.invoke_user_id = 1
pipeline.file_name = "test.md"
pipeline.db_file = SimpleNamespace(
id=85044,
knowledge_id=3124,
create_time=datetime(2026, 7, 14, 10, 0, 0),
update_time=datetime(2026, 7, 14, 10, 1, 0),
updater_id=None,
user_metadata=None,
abstract=None,
)
metadata = pipeline.file_metadata
assert metadata["abstract"] == ""
assert metadata["user_metadata"] == {}
@@ -0,0 +1,45 @@
from types import SimpleNamespace
from unittest.mock import MagicMock
from bisheng.common.constants.vectorstore_metadata import KNOWLEDGE_RAG_METADATA_SCHEMA
from bisheng.worker.knowledge import rebuild_knowledge_worker
def test_rebuild_uses_explicit_milvus_metadata_schema(monkeypatch):
knowledge = SimpleNamespace(id=3124, index_name="knowledge-index", collection_name="knowledge-collection")
embeddings = MagicMock()
embeddings.embed_query.return_value = [0.1, 0.2]
es_client = SimpleNamespace(client=SimpleNamespace(indices=SimpleNamespace(exists=lambda index: True)))
vector_client = MagicMock()
init_milvus = MagicMock(return_value=vector_client)
monkeypatch.setattr(
rebuild_knowledge_worker.KnowledgeRag,
"init_knowledge_es_vectorstore_sync",
lambda knowledge: es_client,
)
monkeypatch.setattr(
rebuild_knowledge_worker.LLMService,
"get_bisheng_knowledge_embedding_sync",
lambda model_id, invoke_user_id: embeddings,
)
monkeypatch.setattr(
rebuild_knowledge_worker.KnowledgeRag,
"init_knowledge_milvus_vectorstore_sync",
init_milvus,
)
result = rebuild_knowledge_worker._rebuild_embeddings(
knowledge=knowledge,
files=[],
new_model_id=884,
invoke_user_id=1,
)
assert result == ([], [])
init_milvus.assert_called_once_with(
invoke_user_id=1,
knowledge=knowledge,
embeddings=embeddings,
metadata_schemas=KNOWLEDGE_RAG_METADATA_SCHEMA,
)