Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
duplicate_document_indexing_task.py95 linesDownload Raw Back to tasks
1import datetime2import logging3import time4 5import click6from celery import shared_task7 8from configs import dify_config9from core.indexing_runner import DocumentIsPausedError, IndexingRunner10from core.rag.index_processor.index_processor_factory import IndexProcessorFactory11from extensions.ext_database import db12from models.dataset import Dataset, Document, DocumentSegment13from services.feature_service import FeatureService14 15 16@shared_task(queue="dataset")17def duplicate_document_indexing_task(dataset_id: str, document_ids: list):18    """19    Async process document20    :param dataset_id:21    :param document_ids:22 23    Usage: duplicate_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 30    # check document limit31    features = FeatureService.get_features(dataset.tenant_id)32    try:33        if features.billing.enabled:34            vector_space = features.vector_space35            count = len(document_ids)36            batch_upload_limit = int(dify_config.BATCH_UPLOAD_LIMIT)37            if count > batch_upload_limit:38                raise ValueError(f"You have reached the batch upload limit of {batch_upload_limit}.")39            if 0 < vector_space.limit <= vector_space.size:40                raise ValueError(41                    "Your total number of documents plus the number of uploads have over the limit of "42                    "your subscription."43                )44    except Exception as e:45        for document_id in document_ids:46            document = (47                db.session.query(Document).filter(Document.id == document_id, Document.dataset_id == dataset_id).first()48            )49            if document:50                document.indexing_status = "error"51                document.error = str(e)52                document.stopped_at = datetime.datetime.utcnow()53                db.session.add(document)54        db.session.commit()55        return56 57    for document_id in document_ids:58        logging.info(click.style("Start process document: {}".format(document_id), fg="green"))59 60        document = (61            db.session.query(Document).filter(Document.id == document_id, Document.dataset_id == dataset_id).first()62        )63 64        if document:65            # clean old data66            index_type = document.doc_form67            index_processor = IndexProcessorFactory(index_type).init_index_processor()68 69            segments = db.session.query(DocumentSegment).filter(DocumentSegment.document_id == document_id).all()70            if segments:71                index_node_ids = [segment.index_node_id for segment in segments]72 73                # delete from vector index74                index_processor.clean(dataset, index_node_ids)75 76                for segment in segments:77                    db.session.delete(segment)78                db.session.commit()79 80            document.indexing_status = "parsing"81            document.processing_started_at = datetime.datetime.utcnow()82            documents.append(document)83            db.session.add(document)84    db.session.commit()85 86    try:87        indexing_runner = IndexingRunner()88        indexing_runner.run(documents)89        end_at = time.perf_counter()90        logging.info(click.style("Processed dataset: {} latency: {}".format(dataset_id, end_at - start_at), fg="green"))91    except DocumentIsPausedError as ex:92        logging.info(click.style(str(ex), fg="yellow"))93    except Exception:94        pass95