Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
retry_document_indexing_task.py92 linesDownload Raw Back to tasks
1import datetime2import logging3import time4 5import click6from celery import shared_task7 8from core.indexing_runner import IndexingRunner9from core.rag.index_processor.index_processor_factory import IndexProcessorFactory10from extensions.ext_database import db11from extensions.ext_redis import redis_client12from models.dataset import Dataset, Document, DocumentSegment13from services.feature_service import FeatureService14 15 16@shared_task(queue="dataset")17def retry_document_indexing_task(dataset_id: str, document_ids: list[str]):18    """19    Async process document20    :param dataset_id:21    :param document_ids:22 23    Usage: retry_document_indexing_task.delay(dataset_id, document_id)24    """25    documents = []26    start_at = time.perf_counter()27 28    dataset = db.session.query(Dataset).filter(Dataset.id == dataset_id).first()29    for document_id in document_ids:30        retry_indexing_cache_key = "document_{}_is_retried".format(document_id)31        # check document limit32        features = FeatureService.get_features(dataset.tenant_id)33        try:34            if features.billing.enabled:35                vector_space = features.vector_space36                if 0 < vector_space.limit <= vector_space.size:37                    raise ValueError(38                        "Your total number of documents plus the number of uploads have over the limit of "39                        "your subscription."40                    )41        except Exception as e:42            document = (43                db.session.query(Document).filter(Document.id == document_id, Document.dataset_id == dataset_id).first()44            )45            if document:46                document.indexing_status = "error"47                document.error = str(e)48                document.stopped_at = datetime.datetime.utcnow()49                db.session.add(document)50                db.session.commit()51            redis_client.delete(retry_indexing_cache_key)52            return53 54        logging.info(click.style("Start retry document: {}".format(document_id), fg="green"))55        document = (56            db.session.query(Document).filter(Document.id == document_id, Document.dataset_id == dataset_id).first()57        )58        try:59            if document:60                # clean old data61                index_processor = IndexProcessorFactory(document.doc_form).init_index_processor()62 63                segments = db.session.query(DocumentSegment).filter(DocumentSegment.document_id == document_id).all()64                if segments:65                    index_node_ids = [segment.index_node_id for segment in segments]66                    # delete from vector index67                    index_processor.clean(dataset, index_node_ids)68 69                    for segment in segments:70                        db.session.delete(segment)71                    db.session.commit()72 73                document.indexing_status = "parsing"74                document.processing_started_at = datetime.datetime.utcnow()75                db.session.add(document)76                db.session.commit()77 78                indexing_runner = IndexingRunner()79                indexing_runner.run([document])80                redis_client.delete(retry_indexing_cache_key)81        except Exception as ex:82            document.indexing_status = "error"83            document.error = str(ex)84            document.stopped_at = datetime.datetime.utcnow()85            db.session.add(document)86            db.session.commit()87            logging.info(click.style(str(ex), fg="yellow"))88            redis_client.delete(retry_indexing_cache_key)89            pass90    end_at = time.perf_counter()91    logging.info(click.style("Retry dataset: {} latency: {}".format(dataset_id, end_at - start_at), fg="green"))92