Commit 9ccaefd2 authored by Yaowentong's avatar Yaowentong

千问修复

豆包修复
电商数据修复
parent fa820537
This diff is collapsed.
......@@ -176,7 +176,7 @@ def batch_worker():
def main():
thread_list = [
# threading.Thread(target=processing_worker, name="processing_worker", daemon=True),
threading.Thread(target=processing_worker, name="processing_worker", daemon=True),
threading.Thread(target=stream_worker, name="stream_worker", daemon=True),
threading.Thread(target=stream_batch_worker, name="stream_batch_worker", daemon=True),
threading.Thread(target=batch_worker, name="batch_worker", daemon=True),
......
......@@ -6,7 +6,7 @@ from aidso_geo.models.eco_data_process import save_eco_data_to_bh
from aidso_geo.utils import robot_utils, bh_utils
from aidso_geo.utils.ai_interface import get_parse_sse_result
from aidso_geo.utils.tos_utils import get_string_from_tos
from loguru import logger
def doubao_mobile_process_original_data(task_data):
file_path = f'geo/{task_data["taskId"]}/{task_data["platform"]}/original.text'
......@@ -37,6 +37,7 @@ def doubao_mobile_process_original_data(task_data):
try:
json_content = json.loads(payload)
except (IndexError, json.JSONDecodeError):
continue
......@@ -56,7 +57,6 @@ def doubao_mobile_process_original_data(task_data):
response_content += json_content.get("content").get("content_block")[0].get("content").get(
"text_block").get("text")
if is_think:
patch_op = json_content.get('patch_op')
if patch_op:
......@@ -77,6 +77,9 @@ def doubao_mobile_process_original_data(task_data):
if response_bool:
response_content += text_block
# if content_block[0].get('is_finish') and content_block[0].get('parent_id'):
# think_bool = True
# response_bool = False
if content_block[0].get('block_type') == 10025:
if think_bool:
......@@ -93,8 +96,8 @@ def doubao_mobile_process_original_data(task_data):
'search_query_result_block').get(
'results'):
url_list.append(i.get('text_card'))
think_bool = False
response_bool = True
# think_bool = False
# response_bool = True
continue
if content_block[0].get('block_type') == 10040:
think_bool = False
......@@ -146,6 +149,7 @@ def doubao_mobile_process_original_data(task_data):
if content_block[0].get('block_type') == 10101 and content_block[0].get('is_finish'):
response_content = ""
if content_block[0].get('block_type') == 10000 and content_block[0].get('is_finish'):
response_bool = True
if content_block[0].get('block_type') == 10000 and content_block[0].get('is_finish'):
response_bool = False
......@@ -301,6 +305,7 @@ def doubao_mobile_process_original_data(task_data):
if dy_eco_list:
save_eco_data_to_bh(task_data,'dyeco',dy_eco_list)
url_list = new_url_list
spider_save_tos.process_and_save_files(file_path, search_keyword, url_list, think_content, response_content,
suggestions, rich_media_block)
return (file_path, search_keyword, url_list, think_content, response_content, suggestions)
......@@ -337,13 +342,16 @@ if __name__ == '__main__':
# file_path = 'geo/c7eb465e-f385-4aa2-89c4-a7cf11897f45/KIMI/1(1).txt'
data_list = bh_utils.query_data(
f"select * from geo_commit_task where taskId = '369df725-ef53-4050-8b94-233a007950c3' and platform = 'DOUBA'")
f"select * from geo_commit_task where platform = 'DOUBA' and insertime>1783008000 and type !='success'")
# # #
# # # # #
def handle_item(i):
task_id = i.get('taskId')
logger.success(f"{task_id}处理完成")
if i.get('comWordsMap'):
i['comWordsMap'] = json.loads(i.get('comWordsMap'))
if i.get('brandWords'):
......
......@@ -89,5 +89,4 @@ def save_eco_data_to_bh(data,eco_type,eco_list):
...
if eco_type == 'dyeco':
eco_result =doubao_android_process_douyin_eco(data,eco_list)
...
bh_utils.insert_data('geo_eco_data',eco_result)
\ No newline at end of file
......@@ -1954,7 +1954,7 @@ def run_data(PAGE_SIZE,MAX_WORKERS):
if __name__ == '__main__':
data_list = bh_utils.query_data(f"select * from geo_commit_task where status = 'ING' ")
data_list = bh_utils.query_data(f"select * from geo_commit_task where status = 'ING'")
# data_list = bh_utils.query_data(query_sql)
# print(data_list)
# # #
......@@ -1972,10 +1972,10 @@ if __name__ == '__main__':
if i.get('productWordsMap'):
i['productWordsMap'] = json.loads(i.get('productWordsMap'))
type_t = i.get('type')
# type_t = 'batch'
type_t = 'batch'
# commit_task(i,'ING')
# return task_send_queue(i,type_t)
return platform_process(i)
return task_send_queue(i,type_t)
# return platform_process(i)
#
if data_list:
......
......@@ -2,7 +2,7 @@ import json
import traceback
from asyncio import as_completed
from concurrent.futures import ThreadPoolExecutor
from loguru import logger
from aidso_geo.models import spider_save_tos
from aidso_geo.models.eco_data_process import save_eco_data_to_bh
from aidso_geo.utils import robot_utils, bh_utils
......@@ -145,7 +145,24 @@ def qianwen_android_process_original_data(task_data):
if messages[0].get('meta_data').get('elements'):
response_content+=messages[0].get('meta_data').get('elements')[0].get('content')
# 引用来源-无下标
if mime_type == 'multi_load/iframe' :
for i in multi_load:
if i.get('type') =='ref_source_inline' and i.get('content').get('status') =='complete':
multi_load_content = i.get('content')
multi_load_source_seq = i.get('source_seq')
multi_load_content_query_list = multi_load_content.get('query_list')
multi_load_content_url_list = multi_load_content.get('list')
if multi_load_content_query_list:
search_keyword.extend(multi_load_content_query_list)
if multi_load_content_url_list:
url_reslt = {"content": {"list": multi_load_content_url_list}}
url_list_batch.append({
"url": [url_reslt],
"source_seq": multi_load_source_seq
})
if mime_type == 'bar/iframe' and status == 'complete':
if meta_data:
sources = meta_data.get('sources')
for sou in sources:
......@@ -348,9 +365,13 @@ if __name__ == '__main__':
# qianwen_android_process_original_data(file_path2)
data_list = bh_utils.query_data(f"select * from geo_commit_task where reqId = '4b285be0-c62c-48f7-9d72-2892511c306c' and platform = 'TYQWA'")
# data_list = bh_utils.query_data(f"select * from geo_commit_task where taskId = '0001e01d-a8ea-405e-acef-93e4f55abbff' and platform = 'TYQWA'")
data_list = bh_utils.query_data(f"select * from geo_commit_task where platform = 'TYQWA' and insertime>1783008000 and type !='success'")
def handle_item(i):
task_id = i.get('taskId')
if not get_string_from_tos(f"geo/{task_id}/TYQWA/quote.txt"):
logger.success(f"{task_id}处理完成")
if i.get('comWordsMap'):
i['comWordsMap'] = json.loads(i.get('comWordsMap'))
if i.get('brandWords'):
......
# -*- coding: utf-8 -*-
import json
import os
......@@ -258,7 +259,6 @@ def douyinai_process_quote(url_list):
def qianwen_android_process_quote(url_list):
quto_list = []
for index, item in enumerate(url_list):
source_seq= item.get('source_seq', '')
for url in item.get('url'):
......@@ -277,6 +277,7 @@ def qianwen_android_process_quote(url_list):
"source_seq": source_seq or url.get('source_seq'),
}
quto_list.append(raw_data)
return quto_list
def yuanbao_process_quote(url_list):
......@@ -621,3 +622,6 @@ def process_and_save_files_ai(file_path, search_keyword, url_list, think_content
for content, file_name in data_config:
save_data_to_tos_ai(target_dir, content, file_name)
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