Underground-Digital/Workflow-Engine
0
1import time2 3import click4 5import app6from configs import dify_config7from core.rag.datasource.vdb.tidb_on_qdrant.tidb_service import TidbService8from extensions.ext_database import db9from models.dataset import TidbAuthBinding10 11 12@app.celery.task(queue="dataset")13def create_tidb_serverless_task():14 click.echo(click.style("Start create tidb serverless task.", fg="green"))15 tidb_serverless_number = dify_config.TIDB_SERVERLESS_NUMBER16 start_at = time.perf_counter()17 while True:18 try:19 # check the number of idle tidb serverless20 idle_tidb_serverless_number = TidbAuthBinding.query.filter(TidbAuthBinding.active == False).count()21 if idle_tidb_serverless_number >= tidb_serverless_number:22 break23 # create tidb serverless24 iterations_per_thread = 2025 create_clusters(iterations_per_thread)26 27 except Exception as e:28 click.echo(click.style(f"Error: {e}", fg="red"))29 break30 31 end_at = time.perf_counter()32 click.echo(click.style("Create tidb serverless task success latency: {}".format(end_at - start_at), fg="green"))33 34 35def create_clusters(batch_size):36 try:37 new_clusters = TidbService.batch_create_tidb_serverless_cluster(38 batch_size,39 dify_config.TIDB_PROJECT_ID,40 dify_config.TIDB_API_URL,41 dify_config.TIDB_IAM_API_URL,42 dify_config.TIDB_PUBLIC_KEY,43 dify_config.TIDB_PRIVATE_KEY,44 dify_config.TIDB_REGION,45 )46 for new_cluster in new_clusters:47 tidb_auth_binding = TidbAuthBinding(48 cluster_id=new_cluster["cluster_id"],49 cluster_name=new_cluster["cluster_name"],50 account=new_cluster["account"],51 password=new_cluster["password"],52 )53 db.session.add(tidb_auth_binding)54 db.session.commit()55 except Exception as e:56 click.echo(click.style(f"Error: {e}", fg="red"))57 