Commit 8a384426 authored by Yaowentong's avatar Yaowentong

舆情追踪修复

parent fbffc4cc
...@@ -9,9 +9,10 @@ from loguru import logger ...@@ -9,9 +9,10 @@ from loguru import logger
import uuid import uuid
from openpyxl.workbook import Workbook from openpyxl.workbook import Workbook
import time import time
import os import os ,sys
from concurrent.futures import ThreadPoolExecutor, as_completed from concurrent.futures import ThreadPoolExecutor, as_completed
BASE_DIR = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
sys.path.append(BASE_DIR)
from aidso_geo.models import spider_save_tos from aidso_geo.models import spider_save_tos
from aidso_geo.utils import bh_utils, tos_utils from aidso_geo.utils import bh_utils, tos_utils
...@@ -1277,10 +1278,7 @@ if __name__ == "__main__": ...@@ -1277,10 +1278,7 @@ if __name__ == "__main__":
"run_task(每天0:10)," "run_task(每天0:10),"
) )
scheduler.start() scheduler.start()
# platform_list = ["DB", "DP", "TXYB", "TYQW"]
# for platform in platform_list:
# init_task('20260610', platform)
# run_once_for_pt('20260610')
......
...@@ -293,7 +293,7 @@ def check_quto(): ...@@ -293,7 +293,7 @@ def check_quto():
else: else:
urls.append(c.get('title')) urls.append(c.get('title'))
url_ids_map = url_utils.generate_numeric_url_id(urls) url_ids_map = url_utils.generate_numeric_url_id_v2(urls)
for c in quote_data: for c in quote_data:
u = c.get('url') u = c.get('url')
......
import json import json
import traceback
from aidso_geo.utils import robot_utils, bh_utils from aidso_geo.utils import robot_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
...@@ -20,7 +21,9 @@ def kimi_process_original_data(file_path): ...@@ -20,7 +21,9 @@ def kimi_process_original_data(file_path):
try: try:
original_content = get_string_from_tos(file_path) original_content = get_string_from_tos(file_path)
original_content_list = json.loads(original_content) original_content_list = json.loads(original_content)
if isinstance(original_content_list, list): if isinstance(original_content_list, list):
for item in original_content_list: for item in original_content_list:
block_data = item.get('block', {}) block_data = item.get('block', {})
...@@ -86,13 +89,13 @@ def kimi_process_original_data(file_path): ...@@ -86,13 +89,13 @@ def kimi_process_original_data(file_path):
if __name__ == '__main__': if __name__ == '__main__':
result = bh_utils.query_data("select * from geo_commit_task where platform = 'KIMI' and insertime >1778428800 order by insertime asc ") # result = bh_utils.query_data("select * from geo_commit_task where platform = 'KIMI' and insertime >1778428800 order by insertime asc ")
# # result = [] # # # result = []
for i in result: # for i in result:
file_path = f"geo/{i.get('taskId')}/KIMI/original.text" # file_path = f"geo/{i.get('taskId')}/KIMI/original.text"
kimi_process_original_data(file_path) # kimi_process_original_data(file_path)
# task = '2b2f460d-3ed6-4e0b-940c-9d1c47e27474' task = '96563dfa-af05-497f-a75c-2b2ae348cf5a'
# file_path = f"geo/{task}/KIMI/original.text" file_path = f"geo/{task}/KIMI/original.text"
# kimi_process_original_data(file_path) kimi_process_original_data(file_path)
# aa = 'https://douchacha-web.tos-cn-beijing.volces.com/geo/2ac65e0ef17642eaa42d440e572f693b/KIMI/original.text' # aa = 'https://douchacha-web.tos-cn-beijing.volces.com/geo/2ac65e0ef17642eaa42d440e572f693b/KIMI/original.text'
...@@ -1623,6 +1623,11 @@ def platform_process(data): ...@@ -1623,6 +1623,11 @@ def platform_process(data):
response_content = None 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): 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)
else: else:
...@@ -1636,7 +1641,6 @@ def platform_process(data): ...@@ -1636,7 +1641,6 @@ def platform_process(data):
return True return True
except Exception as e: except Exception as e:
traceback.print_exc()
logger.exception( logger.exception(
f"{data['reqId']}--{data['platform']}--{data['prompt']}--PROCESS_FAIL" f"{data['reqId']}--{data['platform']}--{data['prompt']}--PROCESS_FAIL"
) )
...@@ -1711,7 +1715,7 @@ def process_success(data): ...@@ -1711,7 +1715,7 @@ def process_success(data):
else: else:
urls.append(c.get('title')) urls.append(c.get('title'))
url_ids_map = url_utils.generate_numeric_url_id(urls) url_ids_map = url_utils.generate_numeric_url_id_v2(urls)
for c in quto_list: for c in quto_list:
u = c.get('url') u = c.get('url')
...@@ -1858,7 +1862,7 @@ if __name__ == '__main__': ...@@ -1858,7 +1862,7 @@ if __name__ == '__main__':
# req_id_sql = ",".join([f"'{req_id}'" for req_id in req_ids]) # req_id_sql = ",".join([f"'{req_id}'" for req_id in req_ids])
# print(req_id_sql) # print(req_id_sql)
# # print(len(req_ids)) # # print(len(req_ids))
query_sql = f"select * from geo_commit_task where reqId in ('e436cdd6-9342-4c67-b894-82bc8dc337a5')" 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)) # print(len(data_list))
# #
...@@ -1880,9 +1884,10 @@ if __name__ == '__main__': ...@@ -1880,9 +1884,10 @@ if __name__ == '__main__':
if i.get('productWordsMap'): if i.get('productWordsMap'):
i['productWordsMap'] = json.loads(i.get('productWordsMap')) i['productWordsMap'] = json.loads(i.get('productWordsMap'))
type_t = i.get('type') type_t = i.get('type')
# type_t = 'batch'
return task_send_queue(i,type_t) # return task_send_queue(i,type_t)
# return platform_process(i) return platform_process(i)
if data_list: if data_list:
with ThreadPoolExecutor(max_workers=30) as executor: with ThreadPoolExecutor(max_workers=30) as executor:
......
...@@ -386,7 +386,7 @@ def save_data_to_tos(target_dir, content, file_name): ...@@ -386,7 +386,7 @@ def save_data_to_tos(target_dir, content, file_name):
urls.append(u) urls.append(u)
else: else:
urls.append(c.get('title')) urls.append(c.get('title'))
url_ids_map = url_utils.generate_numeric_url_id(urls) url_ids_map = url_utils.generate_numeric_url_id_v2(urls)
for c in content: for c in content:
u = c.get('url') u = c.get('url')
...@@ -454,7 +454,7 @@ def save_data_to_tos_ai(target_dir, content, file_name): ...@@ -454,7 +454,7 @@ def save_data_to_tos_ai(target_dir, content, file_name):
urls.append(u) urls.append(u)
else: else:
urls.append(c.get('title')) urls.append(c.get('title'))
url_ids_map = url_utils.generate_numeric_url_id(urls) url_ids_map = url_utils.generate_numeric_url_id_v2(urls)
for c in content: for c in content:
u = c.get('url') u = c.get('url')
......
...@@ -85,7 +85,6 @@ def insert_data(table_name, items) -> bool: ...@@ -85,7 +85,6 @@ def insert_data(table_name, items) -> bool:
print(f"插入数据失败: {str(e)}") print(f"插入数据失败: {str(e)}")
return False return False
finally: finally:
if conn is not None: if conn is not None:
conn.close() conn.close()
......
...@@ -39,9 +39,88 @@ def generate_numeric_url_id(urls): ...@@ -39,9 +39,88 @@ def generate_numeric_url_id(urls):
return result_map return result_map
def generate_numeric_url_id_v2(urls):
"""
根据 URL 列表生成 url -> id 映射:
1. 先查 geo_dim_url 已存在的 URL
2. 不存在的 URL 批量插入
3. 再查一次缺失 URL
4. 返回原始 URL 对应的 id
"""
# 兼容单个字符串
if isinstance(urls, str):
urls = [urls]
if not urls:
return {}
# 清洗 URL,只保留有效字符串,并去重,保持原顺序
url_list = list(dict.fromkeys(
url.strip()
for url in urls
if isinstance(url, str) and url.strip()
))
if not url_list:
return {}
url_id_map = {}
# 1. 先查已存在的 URL
sql = """
SELECT id, url
FROM geo_dim_url
WHERE url IN %(urls)s
"""
rows = query_data(sql, {"urls": tuple(url_list)}) or []
for r in rows:
if r and r.get("url"):
url_id_map[r["url"]] = r.get("id")
# 2. 找出没有查到的 URL
missing_urls = [
url for url in url_list
if url not in url_id_map
]
# 3. 插入缺失 URL
if missing_urls:
insert_list = [
{"url": url}
for url in missing_urls
]
insert_data("geo_dim_url", insert_list)
# 4. 插入后再查一次
rows2 = query_data(sql, {"urls": tuple(missing_urls)}) or []
for r in rows2:
if r and r.get("url"):
url_id_map[r["url"]] = r.get("id")
# 5. 返回结果,保持原始 urls 的映射关系
result_map = {}
for url in urls:
if isinstance(url, str) and url.strip():
clean_url = url.strip()
result_map[url] = url_id_map.get(clean_url)
else:
result_map[url] = None
return result_map
if __name__ == '__main__': if __name__ == '__main__':
urls = ["https://baike.baidu"]
print(generate_numeric_url_id(urls))
urls = ["https://www.iesdouy''''in.com/share/video/760323182770899787'3",'https://beautygol.com/foundation-brands/','https://www.iesdouyin.com/share/video/7593368417316141242','https://www.iesdouyin.com/share/video/7588851834841129403','https://www.iesdouyin.com/share/video/7592518076731383796','https://www.voguescandinavia.com/articles/best-foundation','https://www.iesdouyin.com/share/video/7592931462669650609','https://www.iesdouyin.com/share/video/7611501293928833189','https://www.makeupbrands.org/best-makeup-foundation/','https://m.chinabgao.com/top/brand/103376.html','https://www.womenshealthmag.com/tw/beauty/make-up/g63062750/foundation/','https://www.beauty321.com/post/70204?utm_campaign=post&utm_medium=postinside&utm_source=ETtodayfashion','https://www.elle.com/tw/beauty/news/a70086124/2026-elle-beauty-awards-taiwan/']
print(generate_numeric_url_id_v2(urls))
# {'http: // m.toutiao.com / group / 7598066097035805219 / ': 7421792398252101632, # {'http: // m.toutiao.com / group / 7598066097035805219 / ': 7421792398252101632,
# 'http: // m.toutiao.com / group / 7598067408829514280 / ': 7421792398252101633, # 'http: // m.toutiao.com / group / 7598067408829514280 / ': 7421792398252101633,
# 'https: // www.iesdouyin.com / share / video / 7527608280689904956': 7421792398252101634} # 'https: // www.iesdouyin.com / share / video / 7527608280689904956': 7421792398252101634}
......
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