Underground-Digital/Workflow-Engine
0
1import logging2import time3 4import click5from celery import shared_task6 7from core.rag.index_processor.index_processor_factory import IndexProcessorFactory8from core.rag.models.document import Document9from extensions.ext_database import db10from models.dataset import Dataset, DocumentSegment11from models.dataset import Document as DatasetDocument12 13 14@shared_task(queue="dataset")15def deal_dataset_vector_index_task(dataset_id: str, action: str):16 """17 Async deal dataset from index18 :param dataset_id: dataset_id19 :param action: action20 Usage: deal_dataset_vector_index_task.delay(dataset_id, action)21 """22 logging.info(click.style("Start deal dataset vector index: {}".format(dataset_id), fg="green"))23 start_at = time.perf_counter()24 25 try:26 dataset = Dataset.query.filter_by(id=dataset_id).first()27 28 if not dataset:29 raise Exception("Dataset not found")30 index_type = dataset.doc_form31 index_processor = IndexProcessorFactory(index_type).init_index_processor()32 if action == "remove":33 index_processor.clean(dataset, None, with_keywords=False)34 elif action == "add":35 dataset_documents = (36 db.session.query(DatasetDocument)37 .filter(38 DatasetDocument.dataset_id == dataset_id,39 DatasetDocument.indexing_status == "completed",40 DatasetDocument.enabled == True,41 DatasetDocument.archived == False,42 )43 .all()44 )45 46 if dataset_documents:47 dataset_documents_ids = [doc.id for doc in dataset_documents]48 db.session.query(DatasetDocument).filter(DatasetDocument.id.in_(dataset_documents_ids)).update(49 {"indexing_status": "indexing"}, synchronize_session=False50 )51 db.session.commit()52 53 for dataset_document in dataset_documents:54 try:55 # add from vector index56 segments = (57 db.session.query(DocumentSegment)58 .filter(DocumentSegment.document_id == dataset_document.id, DocumentSegment.enabled == True)59 .order_by(DocumentSegment.position.asc())60 .all()61 )62 if segments:63 documents = []64 for segment in segments:65 document = Document(66 page_content=segment.content,67 metadata={68 "doc_id": segment.index_node_id,69 "doc_hash": segment.index_node_hash,70 "document_id": segment.document_id,71 "dataset_id": segment.dataset_id,72 },73 )74 75 documents.append(document)76 # save vector index77 index_processor.load(dataset, documents, with_keywords=False)78 db.session.query(DatasetDocument).filter(DatasetDocument.id == dataset_document.id).update(79 {"indexing_status": "completed"}, synchronize_session=False80 )81 db.session.commit()82 except Exception as e:83 db.session.query(DatasetDocument).filter(DatasetDocument.id == dataset_document.id).update(84 {"indexing_status": "error", "error": str(e)}, synchronize_session=False85 )86 db.session.commit()87 elif action == "update":88 dataset_documents = (89 db.session.query(DatasetDocument)90 .filter(91 DatasetDocument.dataset_id == dataset_id,92 DatasetDocument.indexing_status == "completed",93 DatasetDocument.enabled == True,94 DatasetDocument.archived == False,95 )96 .all()97 )98 # add new index99 if dataset_documents:100 # update document status101 dataset_documents_ids = [doc.id for doc in dataset_documents]102 db.session.query(DatasetDocument).filter(DatasetDocument.id.in_(dataset_documents_ids)).update(103 {"indexing_status": "indexing"}, synchronize_session=False104 )105 db.session.commit()106 107 # clean index108 index_processor.clean(dataset, None, with_keywords=False)109 110 for dataset_document in dataset_documents:111 # update from vector index112 try:113 segments = (114 db.session.query(DocumentSegment)115 .filter(DocumentSegment.document_id == dataset_document.id, DocumentSegment.enabled == True)116 .order_by(DocumentSegment.position.asc())117 .all()118 )119 if segments:120 documents = []121 for segment in segments:122 document = Document(123 page_content=segment.content,124 metadata={125 "doc_id": segment.index_node_id,126 "doc_hash": segment.index_node_hash,127 "document_id": segment.document_id,128 "dataset_id": segment.dataset_id,129 },130 )131 132 documents.append(document)133 # save vector index134 index_processor.load(dataset, documents, with_keywords=False)135 db.session.query(DatasetDocument).filter(DatasetDocument.id == dataset_document.id).update(136 {"indexing_status": "completed"}, synchronize_session=False137 )138 db.session.commit()139 except Exception as e:140 db.session.query(DatasetDocument).filter(DatasetDocument.id == dataset_document.id).update(141 {"indexing_status": "error", "error": str(e)}, synchronize_session=False142 )143 db.session.commit()144 145 end_at = time.perf_counter()146 logging.info(147 click.style("Deal dataset vector index: {} latency: {}".format(dataset_id, end_at - start_at), fg="green")148 )149 except Exception:150 logging.exception("Deal dataset vector index failed")151 