# 新rabbitmq队列 # 导入 backend/__init__.py 中的所有内容 from .... import * import pika from DCbackend import settings from DCbackend.utils.common import success, fail from backend.models import Knowledgebase, DocumentKbm, File2document, File, Task, TaskSublist, KbmDocumentType def send_to_rabbitmq(queue_name, message): """ 将消息发送到指定的RabbitMQ队列 """ try: connection = pika.BlockingConnection(pika.ConnectionParameters( host=settings.RABBITMQ_HOST, port=settings.RABBITMQ_PORT, credentials=pika.PlainCredentials( settings.RABBITMQ_USER, settings.RABBITMQ_PASSWORD ) )) channel = connection.channel() channel.queue_declare(queue=queue_name, durable=True) channel.basic_publish( exchange='', routing_key=queue_name, body=json.dumps(message), properties=pika.BasicProperties( delivery_mode=2, # 使消息持久化 ) ) connection.close() logger.info(f"消息已发送到队列 {queue_name}") connection = None channel = None return True except Exception as e: logger.error(f"发送消息到RabbitMQ时出错: {str(e)}") return False def analysis(request,KbmService): """ 分析请求并处理RabbitMQ队列中的消息。 Args: request (object): 需要分析的请求对象。 KbmService (object): KbmService对象,用于发送消息到RabbitMQ队列。 Returns: None """ document_id = request.POST.get("document_id") start_page = int(request.POST.get('start_page', 1)) end_page = int(request.POST.get('end_page', -1)) max_tokens = int(request.POST.get('max_tokens', 2048)) if max_tokens == 0: max_tokens = 2048 logger.info(f"开始处理文档 ID: {document_id}") try: document = DocumentKbm.objects.get(id=document_id) if int(document.run) in [1, 5]: # 1: 处理中, 5: 等待处理 logger.info(f"文档 {document_id} 已有队列") return success("文档正在处理中或已经处理完成") # 准备消息 message = { 'document_id': document_id, 'start_page': start_page, 'end_page': end_page, 'max_tokens': max_tokens } # 发送消息到队列 if KbmService.send_to_rabbitmq(settings.RABBITMQ_QUEUE_NAME, message): # 更新文档状态为等待处理 document.run = 5 # 5表示等待处理 document.save() logger.info(f"文档 {document_id} 状态已更新为等待处理") return success("文档已添加到处理队列") else: document.run = 4 document.save() return fail("添加文档到处理队列失败") except DocumentKbm.DoesNotExist: logger.error(f"文档 {document_id} 不存在") document.run = 4 document.save() return fail("文档不存在") except Exception as e: logger.error(f"处理文档 {document_id} 时出错: {str(e)}") document.run = 4 document.save() return fail("处理文档时出错")