Commit ee7620ec authored by Yaowentong's avatar Yaowentong

豆包助手

parent bcf81c1d
import time
from datetime import datetime
from loguru import logger
import requests
import json
import redis
from aidso_geo.models import spider_save_tos
import aidso_geo.utils.bh_utils as bh_utils
from concurrent.futures import ThreadPoolExecutor, as_completed
_redis8_client = None
_redis8_pool = None
def init_redis8():
global _redis8_client, _redis8_pool
if _redis8_client is not None:
return _redis8_client
try:
_redis8_pool = redis.ConnectionPool(
host='redis-cnlfmu7rl14awitrz.redis.ivolces.com',
port=6379,
db=8,
password='aiyingli@@123',
socket_timeout=120,
socket_connect_timeout=5,
decode_responses=True,
retry_on_timeout=True,
health_check_interval=30,
socket_keepalive=True,
max_connections=100,
)
_redis8_client = redis.Redis(connection_pool=_redis8_pool)
_redis8_client.ping()
return _redis8_client
except Exception as e:
logger.error(f"Redis 初始化失败: {e}")
_redis8_client = None
_redis8_pool = None
return None
def get_result(prompt):
url = "https://ark.cn-beijing.volces.com/api/v3/responses"
payload = json.dumps({
"model": "doubao-seed-2-1-pro-260628",
"stream": False,
"tools": [
{
"type": "doubao_app",
"feature": {
"ai_search": {
"type": "enabled",
"role_description": prompt
}
}
}
],
"input": [
{
"type": "message",
"role": "user",
"content": [
{
"type": "input_text",
"text": prompt
}
]
}
]
})
headers = {
'Authorization': 'Bearer ark-7afc3be2-37a8-47fd-9f02-996258a3d305-27da0',
'ark-beta-doubao-app': 'true',
'Content-Type': 'application/json'
}
response_text = requests.request("POST", url, headers=headers, data=payload,timeout=300)
response_text.raise_for_status()
response_json = response_text.json()
result = []
for out in response_json.get('output'):
out_blocks = out.get('blocks')
result.extend(out_blocks)
return result
def get_think_result(prompt):
url = "https://ark.cn-beijing.volces.com/api/v3/responses"
payload = json.dumps({
"model": "doubao-seed-2-1-pro-260628",
"stream": False,
"tools": [
{
"type": "doubao_app",
"feature": {
"reasoning_search": {
"type": "enabled"
}
}
}
],
"input": [
{
"type": "message",
"role": "user",
"content": [
{
"type": "input_text",
"text": prompt
}
]
}
]
})
headers = {
'Authorization': 'Bearer ark-7afc3be2-37a8-47fd-9f02-996258a3d305-27da0',
'ark-beta-doubao-app': 'true',
'Content-Type': 'application/json'
}
response_text = requests.request("POST", url, headers=headers, data=payload,timeout=300)
response_text.raise_for_status()
response_json = response_text.json()
result = []
for out in response_json.get('output'):
out_blocks = out.get('blocks')
result.extend(out_blocks)
return result
def parse_think(think_result):
think_content = ""
response_content = ""
search_word = []
quto = []
suggestions = []
for i in think_result:
if i.get('type') == 'reasoning_text':
think_content += i.get('reasoning_text')
if i.get('type') == 'reasoning_search':
think_content += "\n\n"
think_content += i.get('summary')
think_content += "\n\n"
search_word.extend(i.get('queries'))
results = i.get('results')
quto.extend(results)
if i.get('type') == 'output_text':
response_content += i.get('text')
return (search_word, quto, think_content, response_content,suggestions)
def parse_result(think_result):
think_content = ""
response_content = ""
search_word = []
quto = []
suggestions = []
for i in think_result:
if i.get('type') == 'search':
think_content += i.get('summary')
search_word.extend(i.get('queries'))
results = i.get('results')
quto.extend(results)
if i.get('type') == 'output_text':
response_content += i.get('text')
return (search_word, quto, think_content, response_content,suggestions)
def process_task(tasks):
task_json = json.loads(tasks)
prompt = task_json.get("prompt")
reqId = task_json.get('reqId')
logger.success(f"{reqId}开始处理")
file_path = (
f'geo/{task_json["taskId"]}/'
f'{task_json["platform"]}/original.text'
)
task_json["pt"] = datetime.now().strftime("%Y%m%d")
task_json["source"] = "assistant"
if str(task_json.get("thinkingEnabled")) == "1":
(
search_keyword,
url_list,
think_content,
response_content,
suggestions,
) = parse_think(
get_think_result(prompt)
)
else:
(
search_keyword,
url_list,
think_content,
response_content,
suggestions,
) = parse_result(
get_result(prompt)
)
if not response_content:
raise RuntimeError("接口返回内容为空")
spider_save_tos.process_and_save_files(
file_path,
search_keyword,
url_list,
think_content,
response_content,
suggestions,
)
task_json["status"] = "SUCCESS"
logger.success(f"{reqId}处理完成")
return task_json
if __name__ == "__main__":
redis8 = init_redis8()
while 1:
tasks = redis8.rpop(
"geo:task_commit:JIKE:list",
5,
) or []
if not tasks:
time.sleep(10)
logger.success('当前无任务')
continue
success_tasks = []
success_original_tasks = []
with ThreadPoolExecutor(max_workers=5) as executor:
future_task_map = {
executor.submit(process_task, task): task
for task in tasks
}
for future in as_completed(future_task_map):
original_task = future_task_map[future]
try:
success_tasks.append(future.result())
success_original_tasks.append(original_task)
except Exception as e:
print(f"任务处理失败:{e}")
redis8.lpush(
"geo:task_commit:JIKE:list",
original_task,
)
if success_tasks:
try:
bh_utils.insert_data(
"geo_commit_task",
success_tasks,
)
except Exception:
redis8.lpush(
"geo:task_commit:JIKE:list",
*success_original_tasks,
)
raise
\ No newline at end of file
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment