Commit 1c52bb3a authored by Yaowentong's avatar Yaowentong

豆包新增api重试

豆包历史任务重跑
parent 82b6ccfd
import requests
import json
def get_result(prompt,country):
url = "https://api.scrapeless.com/api/v2/scraper/request"
payload = json.dumps({
"actor": "scraper.chatgpt",
"input": {
"prompt": prompt,
"country": country,
"web_search": True,
"shopping": True
}
})
headers = {
'x-api-token': 'sk_EsiHURJ1OeABPgoRfuyHP6Nt1vRQNwq1Z7uaZ2jNxeYslydgunSmWDoGybh1L5mb',
'Content-Type': 'application/json'
}
response = requests.request("POST", url, headers=headers, data=payload)
print(response.text)
...@@ -278,10 +278,9 @@ if __name__ == '__main__': ...@@ -278,10 +278,9 @@ if __name__ == '__main__':
'XHSA:geo:batch:list', 'XHSA:geo:batch:list',
'geo:task_commit:list'] 'geo:task_commit:list']
# mt:snipaste_v3:only_content # mt:snipaste_v3:only_content
init_redis_1 = init_redis() init_redis_4 = init_redis4()
init_redis_1.delete('DB:geo:batch:list') # print(init_redis_4.delete('mt:snipaste_v3:only_content'))
init_redis_1.delete('DB:geo:stream_batch:list') print(init_redis_4.llen('mt:snipaste_v3:only_content'))
# old_count, new_count = deduplicate_redis_list( # old_count, new_count = deduplicate_redis_list(
# redis_client8, # redis_client8,
# "geo:task_commit:list", # "geo:task_commit:list",
......
...@@ -6,9 +6,72 @@ import json ...@@ -6,9 +6,72 @@ import json
import redis import redis
from aidso_geo.models import spider_save_tos from aidso_geo.models import spider_save_tos
import aidso_geo.utils.bh_utils as bh_utils import aidso_geo.utils.bh_utils as bh_utils
from concurrent.futures import ThreadPoolExecutor, as_completed from concurrent.futures import (
ThreadPoolExecutor,
wait,
FIRST_COMPLETED,
)
MAX_WORKERS = 1
QUEUE_KEY = "geo:task_commit:JIKE:list1"
_redis8_client = None _redis8_client = None
_redis8_pool = None _redis8_pool = None
import threading
import random
REQUEST_INTERVAL = 0.6
request_rate_lock = threading.Lock()
next_request_time = 0.0
def wait_request_slot():
global next_request_time
with request_rate_lock:
now = time.monotonic()
request_time = max(now, next_request_time)
next_request_time = request_time + REQUEST_INTERVAL
wait_seconds = request_time - now
if wait_seconds > 0:
time.sleep(wait_seconds)
def request_ark(url, headers, payload, max_retries=5):
response = None
for attempt in range(max_retries):
wait_request_slot()
response = requests.post(
url,
headers=headers,
data=payload,
timeout=300,
)
if response.status_code != 429:
response.raise_for_status()
return response
retry_after = response.headers.get("Retry-After")
if retry_after:
try:
wait_seconds = float(retry_after)
except ValueError:
wait_seconds = min(2 ** attempt, 16)
else:
wait_seconds = min(2 ** attempt, 16)
wait_seconds += random.uniform(0, 0.5)
logger.warning(
f"触发429,第{attempt + 1}次重试,"
f"等待{wait_seconds:.2f}秒"
)
time.sleep(wait_seconds)
response.raise_for_status()
def init_redis8(): def init_redis8():
global _redis8_client, _redis8_pool global _redis8_client, _redis8_pool
...@@ -76,8 +139,7 @@ def get_result(prompt): ...@@ -76,8 +139,7 @@ def get_result(prompt):
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response_text = requests.request("POST", url, headers=headers, data=payload,timeout=300) response_text = request_ark(url, headers, payload)
response_text.raise_for_status()
response_json = response_text.json() response_json = response_text.json()
result = [] result = []
for out in response_json.get('output'): for out in response_json.get('output'):
...@@ -122,8 +184,7 @@ def get_think_result(prompt): ...@@ -122,8 +184,7 @@ def get_think_result(prompt):
'Content-Type': 'application/json' 'Content-Type': 'application/json'
} }
response_text = requests.request("POST", url, headers=headers, data=payload,timeout=300) response_text = request_ark(url, headers, payload)
response_text.raise_for_status()
response_json = response_text.json() response_json = response_text.json()
result = [] result = []
for out in response_json.get('output'): for out in response_json.get('output'):
...@@ -223,49 +284,54 @@ def process_task(tasks): ...@@ -223,49 +284,54 @@ def process_task(tasks):
if __name__ == "__main__": if __name__ == "__main__":
redis8 = init_redis8() redis8 = init_redis8()
future_task_map = {}
with ThreadPoolExecutor(
max_workers=MAX_WORKERS
) as executor:
while True:
while len(future_task_map) < MAX_WORKERS:
original_task = redis8.rpop(QUEUE_KEY)
if not original_task:
break
while 1: future = executor.submit(
tasks = redis8.rpop( process_task,
"geo:task_commit:JIKE:list", original_task,
1, )
) or [] future_task_map[future] = original_task
if not tasks: if not future_task_map:
time.sleep(10) logger.success("当前无任务")
logger.success('当前无任务') time.sleep(1)
continue continue
success_tasks = [] done, _ = wait(
success_original_tasks = [] tuple(future_task_map),
with ThreadPoolExecutor(max_workers=1) as executor: timeout=1,
future_task_map = { return_when=FIRST_COMPLETED,
executor.submit(process_task, task): task )
for task in tasks
}
for future in as_completed(future_task_map): for future in done:
original_task = future_task_map[future] original_task = future_task_map.pop(future)
try: try:
success_tasks.append(future.result()) success_task = future.result()
success_original_tasks.append(original_task)
except Exception as e:
print(f"任务处理失败:{e}")
redis8.lpush( insert_ok = bh_utils.insert_data(
"geo:task_commit:JIKE:list", "geo_commit_task",
original_task, [success_task],
) )
if success_tasks: if not insert_ok:
try: raise RuntimeError("数据库写入失败")
bh_utils.insert_data(
"geo_commit_task", except Exception as error:
success_tasks, logger.exception(
f"任务处理失败,重新放回队列:{error}"
) )
except Exception:
redis8.lpush( redis8.lpush(
"geo:task_commit:JIKE:list", QUEUE_KEY,
*success_original_tasks, original_task,
) )
\ No newline at end of file
raise
\ No newline at end of file
import time
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,
wait,
FIRST_COMPLETED,
)
MAX_WORKERS = 45
QUEUE_KEY = "mt:snipaste_v3:only_content"
_redis8_client = None
_redis8_pool = None
import threading
import random
REQUEST_INTERVAL = 0.6
request_rate_lock = threading.Lock()
next_request_time = 0.0
def wait_request_slot():
global next_request_time
with request_rate_lock:
now = time.monotonic()
request_time = max(now, next_request_time)
next_request_time = request_time + REQUEST_INTERVAL
wait_seconds = request_time - now
if wait_seconds > 0:
time.sleep(wait_seconds)
def request_ark(url, headers, payload, max_retries=5):
response = None
for attempt in range(max_retries):
wait_request_slot()
response = requests.post(
url,
headers=headers,
data=payload,
timeout=300,
)
if response.status_code != 429:
response.raise_for_status()
return response
retry_after = response.headers.get("Retry-After")
if retry_after:
try:
wait_seconds = float(retry_after)
except ValueError:
wait_seconds = min(2 ** attempt, 16)
else:
wait_seconds = min(2 ** attempt, 16)
wait_seconds += random.uniform(0, 0.5)
logger.warning(
f"触发429,第{attempt + 1}次重试,"
f"等待{wait_seconds:.2f}秒"
)
time.sleep(wait_seconds)
response.raise_for_status()
def init_redis():
try:
return redis.Redis(
host="172.16.0.24",
port=6379,
db=4,
password="aiyingli@@123",
socket_timeout=5,
decode_responses=True,
)
except Exception:
logger.exception("Redis 初始化失败")
return None
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 = request_ark(url, headers, payload)
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 process_task(tasks):
task_json = json.loads(tasks)
prompt = task_json.get("prompt")
reqId = task_json.get('reqId')
pt = task_json.get('pt')
logger.success(f"{reqId}开始处理")
file_path = (
f'geo_snipaste/{pt}/doubao/{reqId}/original.text'
)
(search_keyword,url_list,think_content,response_content,suggestions,) = parse_think(get_think_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["insertime"] = int(time.time())
task_json["platform"] = 'DB'
logger.success(f"{reqId}处理完成")
return task_json
if __name__ == "__main__":
redis8 = init_redis()
future_task_map = {}
with ThreadPoolExecutor(
max_workers=MAX_WORKERS
) as executor:
while True:
while len(future_task_map) < MAX_WORKERS:
original_task = redis8.rpop(QUEUE_KEY)
if not original_task:
break
future = executor.submit(
process_task,
original_task,
)
future_task_map[future] = original_task
if not future_task_map:
logger.success("当前无任务")
time.sleep(1)
continue
done, _ = wait(
tuple(future_task_map),
timeout=1,
return_when=FIRST_COMPLETED,
)
for future in done:
original_task = future_task_map.pop(future)
try:
success_task = future.result()
insert_ok = bh_utils.insert_data(
"geo_third_task",
[success_task],
)
if not insert_ok:
raise RuntimeError("数据库写入失败")
except Exception as error:
logger.exception(
f"任务处理失败,重新放回队列:{error}"
)
redis8.lpush(
QUEUE_KEY,
original_task,
)
\ No newline at end of file
...@@ -3183,4 +3183,13 @@ def run_daily_pipeline_safely( ...@@ -3183,4 +3183,13 @@ def run_daily_pipeline_safely(
if __name__ == "__main__": if __name__ == "__main__":
platform = "DB"
try:
run_daily_pipeline(pt="20260810", platform="DB")
except Exception:
logger.exception(
f"[每日任务 platform={platform} count=1-99] "
f"七阶段任务执行失败"
)
start_scheduler() start_scheduler()
This source diff could not be displayed because it is too large. You can view the blob instead.
...@@ -169,20 +169,17 @@ def get_result_api(): ...@@ -169,20 +169,17 @@ def get_result_api():
def task_commit_api1(): def task_commit_api1():
query_queue = bh_utils.query_data( query_queue = bh_utils.query_data(
""" """
select t1.reqId, SELECT reqId,
t1.prompt, prompt,
t1.taskId, taskId,
t1.platform, platform,
t1.type, type,
t1.insertime, insertime,
t1.status, status,
t1.thinkingEnabled, thinkingEnabled,
t1.channel channel
from (SELECT *
FROM geo_third_task_data FROM geo_third_task_data
where platform = 'DB' ) t1 where platform = 'DB' and status!='SUCCESS'
left join (select * from geo_commit_task where pt > date_format(date_sub(now(), 10), '%Y%m%d')) t2
on t1.reqId = t2.reqId where t2.status ='ING'
""" """
) )
...@@ -263,4 +260,3 @@ if __name__ == "__main__": ...@@ -263,4 +260,3 @@ if __name__ == "__main__":
logger.info("third_task_check 启动") logger.info("third_task_check 启动")
...@@ -5,9 +5,10 @@ from aidso_geo.models import spider_save_tos ...@@ -5,9 +5,10 @@ from aidso_geo.models import spider_save_tos
from aidso_geo.utils import robot_utils, tos_utils, bh_utils from aidso_geo.utils import robot_utils, tos_utils, bh_utils
from aidso_geo.utils.ai_interface import get_parse_sse_result from aidso_geo.utils.ai_interface import get_parse_sse_result
from aidso_geo.utils.tos_utils import get_string_from_tos from aidso_geo.utils.tos_utils import get_string_from_tos
from aidso_geo.config.base_config import init_redis
def doubao_process_original_data(data): def doubao_process_original_data(data):
thinking_enabled = data.get('thinkingEnabled', '0')
file_path = f'geo/{data["taskId"]}/{data["platform"]}/original.text' file_path = f'geo/{data["taskId"]}/{data["platform"]}/original.text'
url_list = "" url_list = ""
think_content = "" think_content = ""
...@@ -21,6 +22,15 @@ def doubao_process_original_data(data): ...@@ -21,6 +22,15 @@ def doubao_process_original_data(data):
try: try:
original_content = get_string_from_tos(file_path) original_content = get_string_from_tos(file_path)
try:
json_assistant = json.loads(original_content)
search_keyword = json_assistant.get('search_word')
url_list = json_assistant.get('quto')
response_content = json_assistant.get('response_content')
think_content = json_assistant.get('think_content')
suggestions = json_assistant.get('suggestions')
except Exception as e:
content_list = original_content.split("\n") content_list = original_content.split("\n")
for i in content_list: for i in content_list:
...@@ -158,7 +168,8 @@ def doubao_process_original_data(data): ...@@ -158,7 +168,8 @@ def doubao_process_original_data(data):
response_content = content_block[0].get("content").get("text_block").get("text") response_content = content_block[0].get("content").get("text_block").get("text")
suggestions = list(set(suggestions)) suggestions = list(set(suggestions))
if str(thinking_enabled) == '1' and not url_list:
return (file_path, [], [], "", "", [])
spider_save_tos.process_and_save_files(file_path, search_keyword, url_list, think_content, response_content, spider_save_tos.process_and_save_files(file_path, search_keyword, url_list, think_content, response_content,
suggestions,rich_media_block) suggestions,rich_media_block)
return (file_path, search_keyword, url_list, think_content, response_content, suggestions) return (file_path, search_keyword, url_list, think_content, response_content, suggestions)
...@@ -190,23 +201,8 @@ def doubao_process_original_data(data): ...@@ -190,23 +201,8 @@ def doubao_process_original_data(data):
if __name__ == '__main__': if __name__ == '__main__':
# task_id_list = [
# '8629188a-46c9-429f-8aa6-d342110326f0', 'ba138c4d-2bb4-45e3-bda3-ecf80553ac89', data_list = bh_utils.query_data(f"select * from geo_commit_task where taskId = '0044cabc-dd1b-493f-a68d-dc7ccdc474b6' and platform = 'DB'")
# 'e178ee36-b9ca-4f1f-802b-83b5e4b820f0', '52efa2d8-616f-48b9-b656-59833fbb91b8',
# '769570cf-c385-4ed7-8e28-f57b5f89d34b', 'cac5bbdb-2145-44a0-9b81-eb20b286d278',
# '7d8ad6ad-c784-407e-9e5f-127c15de5c9e', '6a839541-e3f3-4cd6-a33d-366b3c3297b9',
# '86bd6c56-18d4-41c4-a75a-8d15da5ca589', 'd7e86e4b-521b-4f6f-a34d-ccc7fa629e09',
# '5069cfbb-1150-4bf6-a967-698640c2ce8f', '8bc46734-3469-4681-8a0c-f048169e9486',
# '0480d9c5-0ca9-4c5f-a5c9-dfc4db3b5511', '34cfbfae-cf25-49e6-80a9-199d31813bf2',
# '1d35ff25-760d-43c4-a24f-210d74e82ffe', 'f1d7ca94-7ddf-4548-9b5f-8890d39e7b6c',
# '6fb6c767-54bf-4872-ad60-c52223238d41', '785bf969-4645-4ca7-a9c9-6a5f104a7894',
# 'e52a3c9c-720a-46c5-98f5-efb3e9d8c597', 'ccc66d28-8843-4ea0-a3b0-2700c655260b',
# '3aeee8d2-3086-48c6-a9cb-e29f4a244e34', '67660404-bac4-4bb4-80bb-4b78fd713af0',
# '4b0c64e1-8393-4638-ab98-6b14c579912c', '9217ad06-413a-4161-9be8-6150ffe7b4eb',
# 'e9b27490-6ff0-48da-91a5-3dbbb8494c1d', 'b92b318c71d54c399ed033a722c06a35'
# ]
# for task in task_id_list:
data_list = bh_utils.query_data(f"select * from geo_commit_task where taskId = '9bd06279-57d5-4da7-b8a7-f533fb7365f9' and platform = 'DB'")
# # # # # #
# # # # # #
......
...@@ -1700,10 +1700,11 @@ def platform_process(data): ...@@ -1700,10 +1700,11 @@ def platform_process(data):
commit_task(data, 'PROCESSING') commit_task(data, 'PROCESSING')
platform = data["platform"] platform = data["platform"]
channel = data.get('channel','') channel = data.get('channel','')
thinking_enabled = data.get('thinkingEnabled', '0')
reqId = data.get('reqId') reqId = data.get('reqId')
context_path = f'geo/{data["taskId"]}/{platform}/context.txt' context_path = f'geo/{data["taskId"]}/{platform}/context.txt'
response_content = None response_content = None
url_list= None
# ------------- # -------------
process_func = PLATFORM_PROCESS_MAP.get(platform) process_func = PLATFORM_PROCESS_MAP.get(platform)
...@@ -1711,12 +1712,17 @@ def platform_process(data): ...@@ -1711,12 +1712,17 @@ def platform_process(data):
if check_file_in_tos(context_path): if check_file_in_tos(context_path):
response_content = tos_utils.get_string_from_tos(context_path) response_content = tos_utils.get_string_from_tos(context_path)
elif process_func: elif process_func:
_, _, _, _, response_content, _ = process_func(data) file_path, search_keyword, url_list, think_content, response_content, suggestions = process_func(data)
# ------------- # -------------
if response_content: if response_content:
result_v2(response_content, data) result_v2(response_content, data)
if str(thinking_enabled) == '1' and platform =='DB' and not url_list:
redis_client.lpush("DB:geo:api:list", json.dumps(data))
logger.success(f"{reqId} assistant 处理")
return
else: else:
scheduler(data) scheduler(data)
else: else:
_, _, _, _, response_content, _ = process_func(data) _, _, _, _, response_content, _ = process_func(data)
if response_content: if response_content:
...@@ -2019,12 +2025,11 @@ if __name__ == '__main__': ...@@ -2019,12 +2025,11 @@ if __name__ == '__main__':
# data_list = bh_utils.query_data("select * from geo_commit_task where platform = 'DP' and thinking_enabled = 1 and insertime > 1777564800 and type!='success' order by insertime asc") # data_list = bh_utils.query_data("select * from geo_commit_task where platform = 'DP' and thinking_enabled = 1 and insertime > 1777564800 and type!='success' order by insertime asc")
data_list = bh_utils.query_data("select * from geo_commit_task where status= 'ING' and channel is null") data_list = bh_utils.query_data("select * from geo_commit_task where status !='SUCCESS'")
# data_list = bh_utils.query_data("select * from geo_commit_task where prompt = '飞鹤和君乐宝奶粉的异同点对比' and platform = 'DB'") # data_list = bh_utils.query_data("select * from geo_commit_task where prompt = '飞鹤和君乐宝奶粉的异同点对比' and platform = 'DB'")
# data_list = bh_utils.query_data("select * from geo_commit_task where pt = '20260720' and platform = 'TYQW' and type = 'success' ") # data_list = bh_utils.query_data("select * from geo_commit_task where pt = '20260720' and platform = 'TYQW' and type = 'success' ")
def handle_item(i): def handle_item(i):
logger.success(f"{i.get('taskId')} 处理成功")
i["comWordsMap"] = safe_json_loads(i.get("comWordsMap"), []) i["comWordsMap"] = safe_json_loads(i.get("comWordsMap"), [])
i["brandWords"] = safe_json_loads(i.get("brandWords"), []) i["brandWords"] = safe_json_loads(i.get("brandWords"), [])
i["comWords"] = safe_json_loads(i.get("comWords"), []) i["comWords"] = safe_json_loads(i.get("comWords"), [])
...@@ -2035,7 +2040,7 @@ if __name__ == '__main__': ...@@ -2035,7 +2040,7 @@ if __name__ == '__main__':
# commit_task(i,'ING') # commit_task(i,'ING')
return task_send_queue(i,type_t) return task_send_queue(i,type_t)
# return deepseek_data_process.deepseek_process_original_data(i) # return deepseek_data_process.deepseek_process_original_data(i)
# return platform_process(i) return platform_process(i)
# #
if data_list: if data_list:
with ThreadPoolExecutor(max_workers=50) as executor: with ThreadPoolExecutor(max_workers=50) as executor:
......
This diff is collapsed.
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