123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104 |
- # 新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("处理文档时出错")
|