-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtest_rabbitmq.py
More file actions
118 lines (108 loc) · 5.51 KB
/
Copy pathtest_rabbitmq.py
File metadata and controls
118 lines (108 loc) · 5.51 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
#!/usr/bin/env python3
"""
RabbitMQ队列功能测试脚本
用于测试文档解析队列的发送和接收功能
使用aio_pika库与main.py保持一致
"""
import asyncio
import json
import queue
import os
import aio_pika
from aio_pika import connect_robust, Message, DeliveryMode, exchange
from config import settings
EXCHANGE_NAME = settings.EXCHANGE_NAME
PUB_INVENTION_ROUTING_KEY = settings.PUB_INVENTION_ROUTING_KEY
SUB_INVENTION_ROUTING_KEY = settings.SUB_INVENTION_ROUTING_KEY
INVENTION_QUEUE_NAME = settings.INVENTION_QUEUE_NAME
PUB_TEMPLATE_ROUTING_KEY = settings.PUB_TEMPLATE_ROUTING_KEY
SUB_TEMPLATE_ROUTING_KEY = settings.SUB_TEMPLATE_ROUTING_KEY
TEMPLATE_QUEUE_NAME = settings.TEMPLATE_QUEUE_NAME
RABBITMQ_HOST = settings.RABBITMQ_HOST
RABBITMQ_PORT = settings.RABBITMQ_PORT
RABBITMQ_VIRTUAL_HOST = settings.RABBITMQ_VIRTUAL_HOST
RABBITMQ_USERNAME = settings.RABBITMQ_USERNAME
RABBITMQ_PASSWORD = settings.RABBITMQ_PASSWORD
async def send_test_task(num_tasks, task_type, template_type):
"""发送测试任务到输入队列"""
try:
# 连接RabbitMQ
connection = await connect_robust(
host=RABBITMQ_HOST,
port=RABBITMQ_PORT,
virtualhost=RABBITMQ_VIRTUAL_HOST,
login=RABBITMQ_USERNAME,
password=RABBITMQ_PASSWORD
)
channel = await connection.channel()
if task_type == "1":
routing_key = SUB_INVENTION_ROUTING_KEY
else:
routing_key = SUB_TEMPLATE_ROUTING_KEY
for i in range(num_tasks):
# 创建测试任务
test_task = {
"task_id": f"test_task_00{i}",
"template_type": template_type,
"files": [
# {
# "file_path": "https://s3.kclab.cloud/bucket-78134-shared/test_1.docx?X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Credential=B2sE0fKv1Y1lOpZtge5u%2F20250902%2Fus-east-1%2Fs3%2Faws4_request&X-Amz-Date=20250902T025032Z&X-Amz-Expires=604800&X-Amz-SignedHeaders=host&X-Amz-Signature=ff3189db85d69bdfbe64222b5bc17274acb376df668aba392602ecfe422b224e",
# "file_name": "sample.docx"
# },
# {
# "file_path": "https://s3.kclab.cloud/bucket-78134-shared/test_1.docx?X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Credential=B2sE0fKv1Y1lOpZtge5u%2F20250902%2Fus-east-1%2Fs3%2Faws4_request&X-Amz-Date=20250902T025032Z&X-Amz-Expires=604800&X-Amz-SignedHeaders=host&X-Amz-Signature=ff3189db85d69bdfbe64222b5bc17274acb376df668aba392602ecfe422b224e",
# "file_name": "sample.docx"
# },
{
"file_path": "https://s3.kclab.cloud/bucket-78134-shared/test_1.docx?X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Credential=B2sE0fKv1Y1lOpZtge5u%2F20250921%2Fus-east-1%2Fs3%2Faws4_request&X-Amz-Date=20250921T052834Z&X-Amz-Expires=604800&X-Amz-SignedHeaders=host&X-Amz-Signature=2f4fdece8d889e84d1e8a5056a7cc9bedf3093710529168b00b9f07d17f8c553",
"file_name": "sample.docx"
},
# {
# "file_path": "https://s3.kclab.cloud/bucket-78134-shared/test_1.xlsx?X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Credential=B2sE0fKv1Y1lOpZtge5u%2F20250902%2Fus-east-1%2Fs3%2Faws4_request&X-Amz-Date=20250902T025207Z&X-Amz-Expires=604800&X-Amz-SignedHeaders=host&X-Amz-Signature=b6bf6ecdd26c3ff114a324e6555aad7378b6ac3703fb3295385c684a10863ab5",
# "file_name": "sample.xlsx"
# },
# {
# "file_path": "https://s3.kclab.cloud/bucket-78134-shared/test/1512.08930.pdf?X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Credential=B2sE0fKv1Y1lOpZtge5u%2F20250916%2Fus-east-1%2Fs3%2Faws4_request&X-Amz-Date=20250916T123032Z&X-Amz-Expires=604800&X-Amz-SignedHeaders=host&X-Amz-Signature=4b846c0d2c83844bd3d453c6afb02851e66286066a3ae4ff392a93d20e349411",
# "file_name": "sample.pdf"
# }
]
}
message_body = json.dumps(test_task, ensure_ascii=False)
exchange = await channel.declare_exchange(EXCHANGE_NAME, type="direct", durable=True)
await exchange.publish(
Message(
body=message_body.encode('utf-8')
),
routing_key=routing_key
)
print(f"发送测试任务: {message_body}")
await connection.close()
except Exception as e:
print(f"发送测试任务失败: {str(e)}")
async def main():
"""主函数"""
print("RabbitMQ队列功能测试")
print("=" * 50)
template_type = None
task_type = input("请输入发送任务类型: [1]发明点 [2]模板: ")
if task_type not in ("1", "2"):
print("请输入正确的任务类型")
return
if task_type == "2":
template_type = input("请输入模板类型: [1]化学 [2]电学 [3]机械: ")
if template_type not in ("1", "2", "3"):
print("请输入正确的模板类型")
return
if template_type == "1":
template_type = "化学"
elif template_type == "2":
template_type = "电学"
elif template_type == "3":
template_type = "机械"
num_tasks = int(input("请输入发送任务数量: "))
if not isinstance(num_tasks, int):
print("请输入数字")
return
await send_test_task(num_tasks, task_type, template_type)
if __name__ == "__main__":
asyncio.run(main())