Commit 48730448 authored by Yaowentong's avatar Yaowentong

修复数据

parent 8f636361
......@@ -10,7 +10,9 @@ from copy import deepcopy
from datetime import datetime
from enum import Enum
from aidso_geo.core.down_load_bot import get_req_id
import json
from concurrent.futures import ThreadPoolExecutor, as_completed
from aidso_geo.utils.tos_utils import put_string_to_tos, check_file_in_tos, get_tos_file_size
BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
sys.path.append(BASE_DIR)
......@@ -1620,20 +1622,35 @@ def platform_process(data):
platform = data["platform"]
original_path = f'geo/{data["taskId"]}/{platform}/original.text'
context_path = f'geo/{data["taskId"]}/{platform}/context.txt'
response_content = None
# 冲跑逻辑位 当俩个都有走下面
# process_func = PLATFORM_PROCESS_MAP.get(platform)
# 只有content就走content
# if process_func:
# _, _, _, _, response_content, _ = process_func(original_path)
if check_file_in_tos(context_path):
#——————————————###################
has_original = check_file_in_tos(original_path)
has_context = check_file_in_tos(context_path)
# 重新跑逻辑 当俩个都有走下面
if has_original and has_context:
process_func = PLATFORM_PROCESS_MAP.get(platform)
if process_func:
_, _, _, _, response_content, _ = process_func(original_path)
# 2. 只有 context:直接读 context
elif has_context:
response_content = tos_utils.get_string_from_tos(context_path)
else:
# 3. 只有 original:跑平台解析逻辑
elif has_original:
process_func = PLATFORM_PROCESS_MAP.get(platform)
if process_func:
_, _, _, _, response_content, _ = process_func(original_path)
#——————————————###################
# if check_file_in_tos(context_path):
# response_content = tos_utils.get_string_from_tos(context_path)
# else:
# process_func = PLATFORM_PROCESS_MAP.get(platform)
# if process_func:
# _, _, _, _, response_content, _ = process_func(original_path)
if response_content:
result_v2(response_content, data)
else:
......@@ -1827,70 +1844,46 @@ def task_send_queue(data, queue):
commit_task(data, 'ING')
redis_client.lpush(f"{data['platform']}:geo:{queue}:list", json.dumps(data))
def handle_item(i):
if i.get('comWordsMap'):
i['comWordsMap'] = json.loads(i.get('comWordsMap'))
if i.get('brandWords'):
i['brandWords'] = json.loads(i.get('brandWords'))
if i.get('comWords'):
i['comWords'] = json.loads(i.get('comWords'))
if i.get('keywords'):
i['keywords'] = json.loads(i.get('keywords'))
if i.get('productWordsMap'):
i['productWordsMap'] = json.loads(i.get('productWordsMap'))
type_t = i.get('type')
# type_t = 'batch'
# return task_send_queue(i,type_t)
return platform_process(i)
def run_data(PAGE_SIZE,MAX_WORKERS):
last_insertime = 1781253622
while True:
query_sql = f"""
SELECT *
FROM geo_commit_task
WHERE insertime < {last_insertime} and status ='SUCCESS'
ORDER BY insertime DESC
LIMIT {PAGE_SIZE}
"""
data_list = bh_utils.query_data(query_sql) or []
if not data_list:
logger.success("处理结束")
break
logger.success(
f"本批次查询到 {len(data_list)} 条,last_insertime={last_insertime}"
)
if __name__ == '__main__':
# data = {
# "prompt": "商家出餐但无骑手接单怎么办?",
# "taskId": "aa609471-94f3-4f36-ae0c-9d5e6377e929",
# "brandWords":["美团APP", "美团"],
# "comWords":["高德本地生活", "京东本地生活", "快手本地生活", "抖音生活服务", "饿了么"],
# "reqId": "45d69530-862b-4a59-b145-355d9c2b1003",
# "platform": "DPA",
# "type": "batch",
# "searchEnabled": 1,
# "thinkingEnabled": 0,
# "comWordsMap":[{"brand": "饿了么", "keywords": ["饿了么"]}, {"brand": "抖音生活服务", "keywords": ["抖音生活服务"]}, {"brand": "高德本地生活", "keywords": ["高德本地生活"]}, {"brand": "京东本地生活", "keywords": ["京东本地生活"]}, {"brand": "快手本地生活", "keywords": ["快手本地生活"]}]
# }
# platform_process(data)
import json
from concurrent.futures import ThreadPoolExecutor, as_completed
#
# begin = '2026-05-28'
# end = '2026-06-04'
# req_list = get_req_id(18900000004,begin,end)
# b = "开云集团"
# req_ids = []
# for item in req_list:
# if item.get('brand_name') == b:
# req_id = item.get("req_id")
# req_ids.append(req_id)
# req_id_sql = ",".join([f"'{req_id}'" for req_id in req_ids])
# print(req_id_sql)
# # print(len(req_ids))
query_sql = f"select * from geo_commit_task where reqId in ('2ed850cc-ef78-4558-8987-ee7839fc76e2','fac3f120-4039-4c16-b4a1-0373f99502e3','28e9a2d6-aae7-4ada-986d-959b11ab0442','7a985079-1b69-4e0f-bfad-8249f629dda1','86e717a9-4268-4112-bcf4-8ad3dced1291','674b3816-f63b-40ca-b327-3af2f31b033d','23df97c9-e1d2-42fc-b3a4-23c60c5d707b','9946b538-9a51-4f33-8de7-497f30fc5abb','3a188141-2d38-4101-91c1-16217204bbed','ca1d0ba5-46fc-4e62-8262-c4adc0efb974','702e39f7-0d66-4e0e-a792-284dd4943267','9a2fa801-59e5-4f25-8402-40afc2d7fbd2','e48b3fac-dca5-49af-92f6-e8610d323b1e','a2d0aab5-593e-448d-9ca6-921aef6d79b5','3cfb67ff-fca4-4570-b7d2-ab3fb92d9666','f2762d25-09f5-4353-944a-9881dfe6791a','f8e81976-b9fd-4f7d-87d8-9a19f9398378','f06d8da6-1d83-4703-b4d7-ba048b9044bf','160cb1b4-10e6-4723-9977-2529f0a6bb56','363b2b37-2eb0-4415-a2d1-242965016a59','97131561-0ae6-41f8-968f-1a8805247588','a8cde40c-2f73-45d8-bcdf-4651455377ee','4e4b2bdb-3151-4899-b958-4e1874ceec8d','df1fb0b9-6273-43de-95f1-397949e07911','dc7e9cb4-c46f-48f3-9c5b-f79f618b152c','926b3bfc-6cc2-43d7-a4b2-351e7e0e9919','1566908d-6000-4ac1-a4e6-92cb3575efb3','6b3c1e05-a497-445f-8b2e-1a09bcad8b49','bdbfc0a1-64e5-49b1-9dd2-e2f6c63350db','c74f32f8-7f8f-420c-b310-106945b06a53')"
# print(len(data_list))
#
# data_list = bh_utils.query_data(f"select * from geo_commit_task where status ='ING' and type = 'stream_batch' ")
data_list = bh_utils.query_data(query_sql)
# print(data_list)
# # #
# # #
# # # # #
def handle_item(i):
if i.get('comWordsMap'):
i['comWordsMap'] = json.loads(i.get('comWordsMap'))
if i.get('brandWords'):
i['brandWords'] = json.loads(i.get('brandWords'))
if i.get('comWords'):
i['comWords'] = json.loads(i.get('comWords'))
if i.get('keywords'):
i['keywords'] = json.loads(i.get('keywords'))
if i.get('productWordsMap'):
i['productWordsMap'] = json.loads(i.get('productWordsMap'))
type_t = i.get('type')
# type_t = 'batch'
# return task_send_queue(i,type_t)
return platform_process(i)
if data_list:
with ThreadPoolExecutor(max_workers=30) as executor:
with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor:
futures = [executor.submit(handle_item, i) for i in data_list]
for future in as_completed(futures):
......@@ -1898,3 +1891,46 @@ if __name__ == '__main__':
future.result()
except Exception as e:
logger.exception(f"platform_process 执行异常: {e}")
# 用本批次最后一条的 insertime 推进游标
last_insertime = int(data_list[-1].get("insertime"))
logger.success(f"本批次处理完成,更新 last_insertime={last_insertime}")
if __name__ == '__main__':
# data_list = bh_utils.query_data(f"select * from geo_commit_task where status ='PROCESSING' ")
# # data_list = bh_utils.query_data(query_sql)
# # print(data_list)
# # # #
# # # #
# # # # # #
# def handle_item(i):
# if i.get('comWordsMap'):
# i['comWordsMap'] = json.loads(i.get('comWordsMap'))
# if i.get('brandWords'):
# i['brandWords'] = json.loads(i.get('brandWords'))
# if i.get('comWords'):
# i['comWords'] = json.loads(i.get('comWords'))
# if i.get('keywords'):
# i['keywords'] = json.loads(i.get('keywords'))
# if i.get('productWordsMap'):
# i['productWordsMap'] = json.loads(i.get('productWordsMap'))
# type_t = i.get('type')
# # type_t = 'batch'
#
# # return task_send_queue(i,type_t)
# return platform_process(i)
#
# if data_list:
# with ThreadPoolExecutor(max_workers=30) as executor:
# futures = [executor.submit(handle_item, i) for i in data_list]
#
# for future in as_completed(futures):
# try:
# future.result()
# except Exception as e:
# logger.exception(f"platform_process 执行异常: {e}")
PAGE_SIZE =10000
MAX_WORKERS =30
run_data(PAGE_SIZE,MAX_WORKERS)
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