Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
deal_dataset_vector_index_task.py151 linesDownload Raw Back to tasks
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