Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
clean_unused_datasets_task.py176 linesDownload Raw Back to schedule
1import datetime2import time3 4import click5from sqlalchemy import func6from werkzeug.exceptions import NotFound7 8import app9from configs import dify_config10from core.rag.index_processor.index_processor_factory import IndexProcessorFactory11from extensions.ext_database import db12from extensions.ext_redis import redis_client13from models.dataset import Dataset, DatasetQuery, Document14from services.feature_service import FeatureService15 16 17@app.celery.task(queue="dataset")18def clean_unused_datasets_task():19    click.echo(click.style("Start clean unused datasets indexes.", fg="green"))20    plan_sandbox_clean_day_setting = dify_config.PLAN_SANDBOX_CLEAN_DAY_SETTING21    plan_pro_clean_day_setting = dify_config.PLAN_PRO_CLEAN_DAY_SETTING22    start_at = time.perf_counter()23    plan_sandbox_clean_day = datetime.datetime.now() - datetime.timedelta(days=plan_sandbox_clean_day_setting)24    plan_pro_clean_day = datetime.datetime.now() - datetime.timedelta(days=plan_pro_clean_day_setting)25    page = 126    while True:27        try:28            # Subquery for counting new documents29            document_subquery_new = (30                db.session.query(Document.dataset_id, func.count(Document.id).label("document_count"))31                .filter(32                    Document.indexing_status == "completed",33                    Document.enabled == True,34                    Document.archived == False,35                    Document.updated_at > plan_sandbox_clean_day,36                )37                .group_by(Document.dataset_id)38                .subquery()39            )40 41            # Subquery for counting old documents42            document_subquery_old = (43                db.session.query(Document.dataset_id, func.count(Document.id).label("document_count"))44                .filter(45                    Document.indexing_status == "completed",46                    Document.enabled == True,47                    Document.archived == False,48                    Document.updated_at < plan_sandbox_clean_day,49                )50                .group_by(Document.dataset_id)51                .subquery()52            )53 54            # Main query with join and filter55            datasets = (56                db.session.query(Dataset)57                .outerjoin(document_subquery_new, Dataset.id == document_subquery_new.c.dataset_id)58                .outerjoin(document_subquery_old, Dataset.id == document_subquery_old.c.dataset_id)59                .filter(60                    Dataset.created_at < plan_sandbox_clean_day,61                    func.coalesce(document_subquery_new.c.document_count, 0) == 0,62                    func.coalesce(document_subquery_old.c.document_count, 0) > 0,63                )64                .order_by(Dataset.created_at.desc())65                .paginate(page=page, per_page=50)66            )67 68        except NotFound:69            break70        if datasets.items is None or len(datasets.items) == 0:71            break72        page += 173        for dataset in datasets:74            dataset_query = (75                db.session.query(DatasetQuery)76                .filter(DatasetQuery.created_at > plan_sandbox_clean_day, DatasetQuery.dataset_id == dataset.id)77                .all()78            )79            if not dataset_query or len(dataset_query) == 0:80                try:81                    # remove index82                    index_processor = IndexProcessorFactory(dataset.doc_form).init_index_processor()83                    index_processor.clean(dataset, None)84 85                    # update document86                    update_params = {Document.enabled: False}87 88                    Document.query.filter_by(dataset_id=dataset.id).update(update_params)89                    db.session.commit()90                    click.echo(click.style("Cleaned unused dataset {} from db success!".format(dataset.id), fg="green"))91                except Exception as e:92                    click.echo(93                        click.style("clean dataset index error: {} {}".format(e.__class__.__name__, str(e)), fg="red")94                    )95    page = 196    while True:97        try:98            # Subquery for counting new documents99            document_subquery_new = (100                db.session.query(Document.dataset_id, func.count(Document.id).label("document_count"))101                .filter(102                    Document.indexing_status == "completed",103                    Document.enabled == True,104                    Document.archived == False,105                    Document.updated_at > plan_pro_clean_day,106                )107                .group_by(Document.dataset_id)108                .subquery()109            )110 111            # Subquery for counting old documents112            document_subquery_old = (113                db.session.query(Document.dataset_id, func.count(Document.id).label("document_count"))114                .filter(115                    Document.indexing_status == "completed",116                    Document.enabled == True,117                    Document.archived == False,118                    Document.updated_at < plan_pro_clean_day,119                )120                .group_by(Document.dataset_id)121                .subquery()122            )123 124            # Main query with join and filter125            datasets = (126                db.session.query(Dataset)127                .outerjoin(document_subquery_new, Dataset.id == document_subquery_new.c.dataset_id)128                .outerjoin(document_subquery_old, Dataset.id == document_subquery_old.c.dataset_id)129                .filter(130                    Dataset.created_at < plan_pro_clean_day,131                    func.coalesce(document_subquery_new.c.document_count, 0) == 0,132                    func.coalesce(document_subquery_old.c.document_count, 0) > 0,133                )134                .order_by(Dataset.created_at.desc())135                .paginate(page=page, per_page=50)136            )137 138        except NotFound:139            break140        if datasets.items is None or len(datasets.items) == 0:141            break142        page += 1143        for dataset in datasets:144            dataset_query = (145                db.session.query(DatasetQuery)146                .filter(DatasetQuery.created_at > plan_pro_clean_day, DatasetQuery.dataset_id == dataset.id)147                .all()148            )149            if not dataset_query or len(dataset_query) == 0:150                try:151                    features_cache_key = f"features:{dataset.tenant_id}"152                    plan = redis_client.get(features_cache_key)153                    if plan is None:154                        features = FeatureService.get_features(dataset.tenant_id)155                        redis_client.setex(features_cache_key, 600, features.billing.subscription.plan)156                        plan = features.billing.subscription.plan157                    if plan == "sandbox":158                        # remove index159                        index_processor = IndexProcessorFactory(dataset.doc_form).init_index_processor()160                        index_processor.clean(dataset, None)161 162                        # update document163                        update_params = {Document.enabled: False}164 165                        Document.query.filter_by(dataset_id=dataset.id).update(update_params)166                        db.session.commit()167                        click.echo(168                            click.style("Cleaned unused dataset {} from db success!".format(dataset.id), fg="green")169                        )170                except Exception as e:171                    click.echo(172                        click.style("clean dataset index error: {} {}".format(e.__class__.__name__, str(e)), fg="red")173                    )174    end_at = time.perf_counter()175    click.echo(click.style("Cleaned unused dataset from db success latency: {}".format(end_at - start_at), fg="green"))176