Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
sync_website_document_indexing_task.py89 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 sync_website_document_indexing_task(dataset_id: str, document_id: str):18    """19    Async process document20    :param dataset_id:21    :param document_id:22 23    Usage: sync_website_document_indexing_task.delay(dataset_id, document_id)24    """25    start_at = time.perf_counter()26 27    dataset = db.session.query(Dataset).filter(Dataset.id == dataset_id).first()28 29    sync_indexing_cache_key = "document_{}_is_sync".format(document_id)30    # check document limit31    features = FeatureService.get_features(dataset.tenant_id)32    try:33        if features.billing.enabled:34            vector_space = features.vector_space35            if 0 < vector_space.limit <= vector_space.size:36                raise ValueError(37                    "Your total number of documents plus the number of uploads have over the limit of "38                    "your subscription."39                )40    except Exception as e:41        document = (42            db.session.query(Document).filter(Document.id == document_id, Document.dataset_id == dataset_id).first()43        )44        if document:45            document.indexing_status = "error"46            document.error = str(e)47            document.stopped_at = datetime.datetime.utcnow()48            db.session.add(document)49            db.session.commit()50        redis_client.delete(sync_indexing_cache_key)51        return52 53    logging.info(click.style("Start sync website document: {}".format(document_id), fg="green"))54    document = db.session.query(Document).filter(Document.id == document_id, Document.dataset_id == dataset_id).first()55    try:56        if document:57            # clean old data58            index_processor = IndexProcessorFactory(document.doc_form).init_index_processor()59 60            segments = db.session.query(DocumentSegment).filter(DocumentSegment.document_id == document_id).all()61            if segments:62                index_node_ids = [segment.index_node_id for segment in segments]63                # delete from vector index64                index_processor.clean(dataset, index_node_ids)65 66                for segment in segments:67                    db.session.delete(segment)68                db.session.commit()69 70            document.indexing_status = "parsing"71            document.processing_started_at = datetime.datetime.utcnow()72            db.session.add(document)73            db.session.commit()74 75            indexing_runner = IndexingRunner()76            indexing_runner.run([document])77            redis_client.delete(sync_indexing_cache_key)78    except Exception as ex:79        document.indexing_status = "error"80        document.error = str(ex)81        document.stopped_at = datetime.datetime.utcnow()82        db.session.add(document)83        db.session.commit()84        logging.info(click.style(str(ex), fg="yellow"))85        redis_client.delete(sync_indexing_cache_key)86        pass87    end_at = time.perf_counter()88    logging.info(click.style("Sync document: {} latency: {}".format(document_id, end_at - start_at), fg="green"))89