Commit fbffc4cc authored by Yaowentong's avatar Yaowentong

飞书截图_v2

parent d9c9e149
import json
from datetime import datetime, timedelta
from apscheduler.schedulers.blocking import BlockingScheduler
import openpyxl
import requests
from urllib.parse import urlparse, parse_qs
import redis
from loguru import logger
import uuid
from openpyxl.workbook import Workbook
import time
import os
from concurrent.futures import ThreadPoolExecutor, as_completed
from aidso_geo.models import spider_save_tos
from aidso_geo.utils import bh_utils, tos_utils
APP_ID = "cli_aaac8fa146fadbd5"
APP_SECRET = "N5omKg6d3sHi3JEIopQ8lcDAZKOiaNQ4"
SHEET_URL = "https://s12is4u3s19.feishu.cn/sheets/THiJsNAukhTuwjt47kgc0TiZneb"
SHEET_ID = "THiJsNAukhTuwjt47kgc0TiZneb"
POLL_SECONDS = 180
EXCEL_DIR = "/Users/yaowentong/Desktop"
INIT_TASK_LOCK_EXPIRE = 3 * 24 * 3600 # 3天:防止 key 堆积
container_id = 'oc_60a61d199f4a81b673f01b489b40509b'
def parse_spreadsheet_token(sheet_url: str) -> str:
parsed = urlparse(sheet_url)
parts = parsed.path.strip("/").split("/")
if len(parts) >= 2 and parts[0] == "sheets":
return parts[1]
raise ValueError(f"无法解析 spreadsheet_token: {sheet_url}")
def request_json(method, url, **kwargs):
resp = requests.request(method, url, timeout=30, **kwargs)
try:
data = resp.json()
except Exception:
raise RuntimeError(f"返回不是 JSON: {resp.text}")
if resp.status_code >= 400:
raise RuntimeError(f"HTTP 请求失败: status={resp.status_code}, body={data}")
if data.get("code") != 0:
raise RuntimeError(f"飞书业务错误: {data}")
return data
def get_tenant_access_token(app_id: str, app_secret: str) -> str:
url = "https://open.feishu.cn/open-apis/auth/v3/tenant_access_token/internal"
data = request_json(
"POST",
url,
json={
"app_id": app_id,
"app_secret": app_secret
}
)
return data["tenant_access_token"]
def get_sheet_meta(spreadsheet_token: str, tenant_access_token: str) -> dict:
url = f"https://open.feishu.cn/open-apis/sheets/v3/spreadsheets/{spreadsheet_token}/sheets/query"
headers = {
"Authorization": f"Bearer {tenant_access_token}",
"Content-Type": "application/json; charset=utf-8"
}
return request_json("GET", url, headers=headers)
def get_first_sheet_id(spreadsheet_token: str, tenant_access_token: str) -> str:
data = get_sheet_meta(spreadsheet_token, tenant_access_token)
sheets = data.get("data", {}).get("sheets", [])
if not sheets:
raise RuntimeError(f"没有获取到 sheets: {data}")
first_sheet = sheets[0]
sheet_id = first_sheet.get("sheet_id")
if not sheet_id:
raise RuntimeError(f"没有 sheet_id: {first_sheet}")
return sheet_id
def batch_get_values(spreadsheet_token: str, tenant_access_token: str, ranges: list):
url = f"https://open.feishu.cn/open-apis/sheets/v2/spreadsheets/{spreadsheet_token}/values_batch_get"
headers = {
"Authorization": f"Bearer {tenant_access_token}",
"Content-Type": "application/json; charset=utf-8"
}
params = []
for r in ranges:
params.append(("ranges", r))
return request_json("GET", url, headers=headers, params=params)
def get_not_null_values(values):
"""
把飞书返回的二维数组,转成一维数组,并过滤 None
"""
result = []
for row in values:
for cell in row:
if cell is not None:
result.append(cell)
return result
def init_redis4():
try:
redis_client = redis.Redis(
host="172.16.0.24",
port=6379,
db=4,
password="aiyingli@@123",
socket_timeout=5,
decode_responses=True
)
return redis_client
except Exception as e:
return None
def init_task(pt_for_tomorrow: str, platform):
spreadsheet_token = parse_spreadsheet_token(SHEET_URL)
tenant_access_token = get_tenant_access_token(APP_ID, APP_SECRET)
sheet_id = get_first_sheet_id(spreadsheet_token, tenant_access_token)
platform_map = {
'DB': [
f"{sheet_id}!A2:A10000"
],
'DP': [
f"{sheet_id}!B2:B10000"
],
'TXYB': [
f"{sheet_id}!C2:C10000"
],
'TYQW': [
f"{sheet_id}!D2:D10000"
],
}
result = batch_get_values(
spreadsheet_token=spreadsheet_token,
tenant_access_token=tenant_access_token,
ranges=platform_map[platform]
)
value_ranges = result.get("data", {}).get("valueRanges", [])
data_list = []
for value_range in value_ranges:
values = value_range.get("values", [])
data = get_not_null_values(values)
for cn in range(3):
for index,p in enumerate(data):
req_id = str(uuid.uuid4())
data_list.append({
"reqId": req_id,
"prompt": p,
"platform": platform,
"cn": cn,
"pt": pt_for_tomorrow,
"rank": index,
"channel":'MT'
})
bh_utils.insert_data("geo_feishu_snipaste", data_list)
logger.success(f"[init_task] 已写入任务: pt={pt_for_tomorrow},platform={platform}, rows={len(data_list)}")
def init_task_once(pt_for_tomorrow: str, platform):
"""
只对某个 pt 生成一次任务(避免 scheduler 线程重复触发或重启重复写)
"""
r = init_redis4()
if not r:
# redis 不可用时,退化为直接执行(可能重复写入,至少不影响流程)
init_task(pt_for_tomorrow, platform)
return
lock_key = f"mt:init_task_done:{pt_for_tomorrow}:{platform}"
ok = r.set(lock_key, "1", nx=True, ex=INIT_TASK_LOCK_EXPIRE)
if ok:
init_task(pt_for_tomorrow, platform)
else:
logger.success(f"[init_task_once] 已生成过 pt={pt_for_tomorrow},跳过")
def init_task_scheduler():
tomorrow = datetime.now().date() + timedelta(days=1)
pt_tomorrow = tomorrow.strftime("%Y%m%d")
plate_form =['DB','DP','TXYB','TYQW']
for i in plate_form:
init_task_once(pt_tomorrow,i)
def wait_mt_tasks_done(redis_client,platform):
while True:
try:
size = redis_client.scard(f"mt:feishu_snipaste:{platform}")
logger.success(f"platform:{platform}size:{size}")
except Exception as e:
logger.success(f"redis 读取 mt:feishu_snipaste:{platform} 失败: {e}")
size = -1
if size == 0:
return
time.sleep(POLL_SECONDS)
def doubao_process_original_data(file_path, original_content):
url_list = ""
think_content = ""
response_content = ""
search_keyword = []
suggestions = []
is_think = False
rich_media_block = []
think_bool = False
response_bool = False
content_list = original_content.split("\n")
for i in content_list:
if i != "":
if not i.startswith("data:"):
continue
payload = i[len("data:"):].lstrip()
if not payload:
continue
try:
json_content = json.loads(payload)
except (IndexError, json.JSONDecodeError):
continue
if json_content.get('event_type') == 2001:
even_data = json.loads(json_content.get('event_data'))
message_data = even_data.get('message')
if even_data.get('tts_content') is not None:
response_content = even_data.get('tts_content')
if message_data.get('content_type') == 2007:
for i in json.loads(message_data.get('content')).get("search_result").get("video_card").get(
"card_list"):
rich_media_block.append(i)
if message_data.get('content_type') == 10040 and message_data.get('is_finish') is None:
think_bool = True
continue
if message_data.get('content_type') == 10040 and message_data.get('is_finish') == True:
think_bool = False
continue
if think_bool:
if json.loads(message_data.get('content')).get('text') is not None:
think_content += json.loads(message_data.get('content')).get('text')
content_json = json.loads(message_data.get('content'))
if message_data.get('content_type') == 10025 and content_json.get('results') is not None:
think_content += "\n\n"
think_content += "**搜索"
think_content += str(len(content_json.get('queries')))
think_content += "个关键词,参考"
think_content += str(len(content_json.get('results')))
think_content += "篇文章**"
think_content += "\n\n"
if message_data.get('content_type') == 10025:
content_json = json.loads(message_data.get('content'))
if content_json.get('queries') is not None and content_json.get('results') is not None:
search_keyword = search_keyword + json.loads(message_data.get('content')).get('queries')
if content_json.get('scene') == 2:
url_list = content_json.get('results')
if message_data.get('content_type') == 2002:
suggestions = suggestions + json.loads(message_data.get('content')).get('suggestions')
else:
if json_content.get("patch_op"):
if json_content.get('patch_op')[0].get("patch_object") == 111:
if json_content.get('patch_op')[0].get("patch_value").get("tts_content"):
response_content += json_content.get('patch_op')[0].get("patch_value").get("tts_content")
if json_content.get('patch_op')[0].get("patch_object") == 1:
if json_content.get('patch_op')[0].get("patch_value").get("content_block")[0].get(
"content").get("search_query_result_block"):
search_keyword = json_content.get('patch_op')[0].get("patch_value").get("content_block")[
0].get("content").get("search_query_result_block").get("queries")
url_list = json_content.get('patch_op')[0].get("patch_value").get("content_block")[0].get(
"content").get("search_query_result_block").get("results")
if is_think == True and \
json_content.get('patch_op')[0].get("patch_value").get("content_block")[0].get(
"content").get("search_query_result_block").get("results") is not None:
think_content += "\n\n"
think_content += "**"
think_content += \
json_content.get('patch_op')[0].get("patch_value").get("content_block")[0].get(
"content").get("search_query_result_block").get("summary")
think_content += "**"
think_content += "\n\n"
if json_content.get('patch_op')[0].get("patch_value").get("content_block")[0].get(
"block_type") == 10000 and \
json_content.get('patch_op')[0].get("patch_value").get("content_block")[0].get(
"parent_id") and len(json_content.get('patch_op')) > 1:
is_think = True
if json_content.get('patch_op')[0].get("patch_value").get("content_block")[0].get(
"content").get('text_block').get("text"):
think_content += \
json_content.get('patch_op')[0].get("patch_value").get("content_block")[0].get(
"content").get("text_block").get("text")
continue
if json_content.get('patch_op')[0].get("patch_value").get("content_block")[0].get(
"block_type") == 10040:
is_think = False
continue
if json_content.get('patch_op')[0].get("patch_object") == 50:
for sug in json.loads(
json_content.get('patch_op')[0].get("patch_value").get("ext").get("sp_v2")):
suggestions.append(sug.get("content"))
if is_think:
if json_content.get("text"):
think_content += json_content.get("text")
else:
if json_content.get("content"):
content_block = json_content.get("content").get('content_block')
if content_block:
if len(content_block) > 0:
if content_block[0].get('block_type') == 10050:
for i in content_block[0].get('content').get('rich_media_block').get('creations'):
rich_media_block.append(i.get('video'))
if content_block[0].get('block_type') == 10000:
if 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))
spider_save_tos.process_and_save_files(file_path, search_keyword, url_list, think_content, response_content,
suggestions)
return (file_path, search_keyword, url_list, think_content, response_content, suggestions)
def deepseek_process_original_data(file_path, original_content):
url_list = ""
think_content = ""
response_content = ""
search_keyword = []
is_think = False
think_bool = False
suggestions = []
response_bool = False
# 按空行分割内容,过滤空字符串
content_list = [item for item in original_content.split("\n\n") if item]
for item in content_list:
if item.startswith("event: "):
# 可根据需要补充event数据处理逻辑
continue
# 处理data类型数据
if item.startswith("data: "):
try:
# 提取并解析JSON数据
data_str = item.split("data: ")[1]
json_data = json.loads(data_str)
except (IndexError, json.JSONDecodeError):
continue # 跳过格式错误的数据
if isinstance(json_data.get('v'),dict):
if json_data.get('v').get('response').get('thinking_enabled') == True:
is_think=True
if json_data.get('v').get('response').get('search_enabled') == True:
if len(json_data.get('v').get('response').get('fragments'))>0:
query_list = json_data.get('v').get('response').get('fragments')[0].get('queries')
if query_list:
result = [item.get('query', '') for item in query_list]
search_keyword.extend(result)
if is_think:
if json_data.get('p') == 'response/fragments' and json_data.get('v')[0].get('type') == 'SEARCH':
search_keyword.append(json_data.get('v')[0].get('queries')[0].get('query'))
if json_data.get('p') == 'response/fragments/0/results' or json_data.get(
'p') == 'response/fragments/-1/results':
url_list = json_data.get('v')
if json_data.get('p') =='response/fragments' and json_data.get('v')[0].get('type') == 'THINK':
think_bool = True
think_content+=json_data.get('v')[0].get('content')
continue
if json_data.get('p') == 'response' :
if isinstance(json_data.get('v')[0].get('v'),list):
if json_data.get('v')[0].get('v')[0].get('type') == 'TOOL_SEARCH':
for q in json_data.get('v')[0].get('v')[0].get('queries'):
search_keyword.append(q.get('query'))
if json_data.get('v')[1].get('p') == 'fragments' or json_data.get('v')[1].get('p') == 'response/fragments':
if json_data.get('v')[1].get('v')[0].get('type') == 'THINK':
think_content+=json_data.get('v')[1].get('v')[0].get('content')
think_bool = True
continue
if json_data.get('v')[0].get('p') == 'fragments' or json_data.get('v')[0].get('p') == 'response/fragments':
if json_data.get('v')[0].get('v')[0].get('type') == 'THINK':
think_content+=json_data.get('v')[0].get('v')[0].get('content')
think_bool = True
continue
if json_data.get('v')[0].get('v') == 'FINISHED':
response_bool = False
if json_data.get('p') == 'response/fragments' and json_data.get('v')[0].get('type') == 'RESPONSE':
think_bool = False
response_bool = True
response_content += json_data.get('v')[0].get('content')
continue
if json_data.get('p') == 'response/fragments/1/elapsed_secs':
think_bool = False
if json_data.get('p') == 'response/fragments/-1/elapsed_secs':
think_bool = False
if json_data.get('p') == 'response/fragments':
response_bool = False
if response_bool:
if isinstance(json_data.get('v'), str):
response_content += json_data.get('v')
if think_bool:
if isinstance(json_data.get('v'), str):
think_content+=json_data.get('v')
else:
if json_data.get('p') == 'response/fragments/-1/content':
response_bool = True
if json_data.get('p') == 'response/fragments' and json_data.get('v')[0].get('type') == 'SEARCH':
search_keyword.append(json_data.get('v')[0].get('queries')[0].get('query'))
if json_data.get('p') == 'response/fragments/0/results' or json_data.get('p') == 'response/fragments/-1/results':
url_list = json_data.get('v')
if json_data.get('p') == 'response/fragments' and json_data.get('v')[0].get('type') == 'RESPONSE':
response_content+=json_data.get('v')[0].get('content')
response_bool = True
continue
if json_data.get('p') == 'response':
if json_data.get('v')[1].get('p') == 'fragments':
if json_data.get('v')[1].get('v')[0].get('type') == 'RESPONSE':
response_content+=json_data.get('v')[1].get('v')[0].get('content')
response_bool = True
continue
if json_data.get('v')[0].get('p') == 'fragments':
if json_data.get('v')[0].get('v')[0].get('type') == 'RESPONSE':
response_content+=json_data.get('v')[0].get('v')[0].get('content')
response_bool = True
continue
if json_data.get('v')[0].get('v') == 'FINISHED':
response_bool = False
if json_data.get('p') == 'response/fragments' and json_data.get('v')[0].get('type') == 'TIP':
response_bool= False
continue
if json_data.get('v') == 'FINISHED':
response_bool = False
if response_bool:
if isinstance(json_data.get('v'),str):
response_content+=json_data.get('v')
spider_save_tos.process_and_save_files(file_path, search_keyword, url_list, think_content, response_content,
suggestions)
return (file_path, search_keyword, url_list, think_content, response_content, suggestions)
def yuanbao_process_original_data(file_path, original_content):
url_list = []
think_content = ""
response_process = ""
response_content = ""
search_keyword = []
suggestions = []
is_think = False
think_bool = False
response_bool = False
rich_media_block =[]
content_list = original_content.split("\n\n")
for i in content_list:
if i.startswith("data: "):
try:
# 提取并解析JSON数据
data_str = i.split("data: ")[1]
json_data = json.loads(data_str)
except (IndexError, json.JSONDecodeError):
continue # 跳过格式错误的数据
if json_data.get('type') == 'searchGuid':
url_list = json_data.get('docs')
if json_data.get('type') == 'think':
think_content+=json_data.get('content')
if json_data.get('type') == 'text' and json_data.get('msg') is not None:
response_content+=json_data.get('msg')
if json_data.get('type') == 'deepSearch':
if json_data.get('contents'):
if json_data.get('contents')[0].get('msg'):
think_content += json_data.get('contents')[0].get('msg')
if not response_process:
response_content += json_data.get('contents')[0].get('msg')
if json_data.get('type') == 'image':
response_content = "生成了图片"
spider_save_tos.process_and_save_files(file_path, search_keyword, url_list, think_content, response_content,
suggestions)
return (file_path, search_keyword, url_list, think_content, response_content, suggestions)
def tongyi_process_original_data(file_path, original_content):
url_list = []
url_list_batch = []
think_content = ""
response_content = ""
search_keyword = []
suggestions = []
is_think = False
think_bool = False
response_bool = False
result = []
content_list = original_content.split("\n")
for i in content_list:
if i.startswith("data:"):
try:
# 提取并解析JSON数据
data_str = i.split("data:")[1]
json_data = json.loads(data_str)
except (IndexError, json.JSONDecodeError):
continue
#
if isinstance(json_data, dict):
if json_data.get('msgStatus'):
if json_data.get("incremental") == False:
if json_data.get('msgStatus') == 'finished':
if json_data.get('contents'):
for cn in json_data.get('contents'):
if cn.get('contentType') == 'plugin':
pluginResult = json.loads(cn.get('content')).get('pluginResult')
if pluginResult:
try:
if isinstance(json.loads(pluginResult), list):
if len(json.loads(pluginResult)) >= 2:
url_list = json.loads(pluginResult)[1].get(
'search_results')
if isinstance(json.loads(pluginResult), dict):
url_list = json.loads(pluginResult).get(
'links')
except Exception as e:
...
if cn.get('contentType') == 'think':
think_content = json.loads(cn.get('content')).get('content')
if cn.get('contentType') == 'text':
response_content = cn.get('content')
if json_data.get("data"):
if json_data.get("data").get('status'):
if json_data.get("data").get('status') == 'complete':
messages = json_data.get('data').get('messages')
for cn in messages:
if cn.get('mime_type') == 'bar/iframe' and cn.get('status') == 'complete':
url_list_batch = cn.get('meta_data').get('sources')[0].get('content').get(
'list')
if cn.get('mime_type') == 'multi_load/iframe' and cn.get('status') == 'complete':
response_content = cn.get('content')
response_content = response_content.replace("[(deep_think)]", "")
multi_load = cn.get('meta_data').get('multi_load')
if multi_load:
for mu in multi_load:
if mu.get('type') == 'deep_think':
if mu.get('content').get('status') == 'complete':
think_content = mu.get('content').get('think_content')
if json_data.get("data").get('messages'):
messages = json_data.get("data").get('messages')
for ms in messages:
if ms.get('mime_type') == 'multi_load/iframe' and ms.get('status') == 'complete':
response_content = ms.get('content')
response_content = response_content.replace("[(deep_think)]", "")
multi_load = ms.get('meta_data').get('multi_load')
if multi_load:
for mu in multi_load:
if mu.get('type') == 'deep_think':
if mu.get('content').get('status') == 'complete':
think_content = mu.get('content').get('think_content')
if ms.get('mime_type') == 'bar/iframe' and ms.get('status') == 'complete':
url_list_batch = ms.get('meta_data').get('sources')[0].get('content').get('list')
if ms.get('mime_type') == 'paa/iframe' and ms.get('status') == 'complete':
paas = ms.get('meta_data').get('paas')
if paas:
for pa in ms.get('meta_data').get('paas'):
suggestions.append(pa.get('show_text'))
if url_list_batch:
for url in url_list_batch:
url_list.append(
{
"url": url.get('url', ''),
"title": url.get('title', ''),
"snippet": url.get('summary', ''),
"host_name": url.get('name', ''),
"host_logo": url.get('icon', ''),
"time": url.get('publish_time', ''),
}
)
spider_save_tos.process_and_save_files(file_path, search_keyword, url_list, think_content, response_content,
suggestions)
return (file_path, search_keyword, url_list, think_content, response_content, suggestions)
PLATFORM_PROCESS_MAP = {
"DB": doubao_process_original_data,
"DP": deepseek_process_original_data,
"TXYB": yuanbao_process_original_data,
"TYQW": tongyi_process_original_data,
}
def timestamp_to_datetime(timestamp, fmt="%Y-%m-%d %H:%M:%S"):
timestamp = int(timestamp)
local_time = datetime.fromtimestamp(timestamp)
return local_time.strftime(fmt)
def get_context(image_url):
url = "https://ark.cn-beijing.volces.com/api/v3/responses"
payload = json.dumps({
"model": "doubao-seed-2-0-lite-260428",
"input": [
{
"role": "user",
"content": [
{
"type": "input_image",
"image_url": image_url
},
{
"type": "input_text",
"text": "提取图中文字"
}
]
}
]
})
headers = {
'Authorization': 'Bearer fcc424e5-58af-494d-9683-5787413a26c9',
'Content-Type': 'application/json'
}
try:
response = requests.request("POST", url, headers=headers, data=payload)
data = response.json()
return next(
(
content.get("text", "")
for item in data.get("output", [])
if item.get("type") == "message"
for content in item.get("content", [])
if content.get("type") == "output_text"
),
""
)
except Exception as e:
return ""
class ExcelWriter:
"""
Excel写入工具类,专门用于写入taskID和response两列数据
"""
def __init__(self, file_path):
"""
初始化Excel写入工具
:param file_path: Excel文件保存路径(如: './output.xlsx')
"""
self.file_path = file_path
self.workbook = None
self.worksheet = None
# 初始化工作簿和工作表
self._init_workbook()
def _init_workbook(self):
"""初始化工作簿和工作表,若文件已存在则打开,不存在则新建"""
# 检查文件是否存在
if os.path.exists(self.file_path):
self.workbook = openpyxl.load_workbook(self.file_path)
# 取第一个工作表
self.worksheet = self.workbook.active
# 检查表头是否存在,不存在则添加
if self.worksheet.cell(row=1, column=1).value != 'taskID' or \
self.worksheet.cell(row=1, column=2).value != 'response':
# 在第一行插入表头
self.worksheet.insert_rows(1)
self.worksheet.cell(row=1, column=1, value='taskID')
self.worksheet.cell(row=1, column=2, value='response')
else:
# 新建工作簿
self.workbook = Workbook()
self.worksheet = self.workbook.active
# 设置表头
self.worksheet.cell(row=1, column=1, value='问题')
self.worksheet.cell(row=1, column=2, value='平台')
self.worksheet.cell(row=1, column=3, value='回答')
self.worksheet.cell(row=1, column=4, value='思考过程')
self.worksheet.cell(row=1, column=5, value='引用来源')
self.worksheet.cell(row=1, column=6, value='分享链接')
self.worksheet.cell(row=1, column=7, value='截图')
self.worksheet.cell(row=1, column=8, value='查询时间')
self.worksheet.cell(row=1, column=9, value='ocr结果')
self.worksheet.cell(row=1, column=10, value='排名')
def write_batch_rows(self, data_list):
if not isinstance(data_list, list) or len(data_list) == 0:
raise ValueError("data_list必须是非空的列表")
# "taskId": i.get("taskId"),
# "prompt": i.get("prompt"),
# "platform": i.get("platform"),
# "insertime": i.get("insertime"),
# 找到下一个空行
next_row = self.worksheet.max_row + 1
# 批量写入数据
for idx, (prompt, platform, content, think, quote,share, snipaste, insertime,ocr,rank) in enumerate(data_list):
self.worksheet.cell(row=next_row + idx, column=1, value=prompt)
self.worksheet.cell(row=next_row + idx, column=2, value=platform)
self.worksheet.cell(row=next_row + idx, column=3, value=content)
self.worksheet.cell(row=next_row + idx, column=4, value=think)
self.worksheet.cell(row=next_row + idx, column=5, value=quote)
self.worksheet.cell(row=next_row + idx, column=6, value=share)
self.worksheet.cell(row=next_row + idx, column=7, value=snipaste)
self.worksheet.cell(row=next_row + idx, column=8, value=insertime)
self.worksheet.cell(row=next_row + idx, column=9, value=ocr)
self.worksheet.cell(row=next_row + idx, column=10, value=rank)
# 保存文件
self.workbook.save(self.file_path)
def close(self):
"""关闭工作簿,释放资源"""
if self.workbook:
self.workbook.close()
def to_excel2(pt, cn,platform):
query_list_log = bh_utils.query_data(
f"SELECT DISTINCT reqId FROM geo_feishu_snipaste where cn = {cn} and pt = {pt} and platform = '{platform}'"
)
uuid_list = []
for i in query_list_log:
uuid_list.append(i.get('reqId'))
uuid_str = ",".join([f"'{uuid}'" for uuid in uuid_list])
query_list = bh_utils.query_data(
f"SELECT DISTINCT reqId,prompt, platform,insertime,rank FROM geo_third_task WHERE reqId IN ({uuid_str}) order by rank asc"
)
task_list = []
for i in query_list:
task_list.append(
{
"reqId": i.get("reqId"),
"prompt": i.get("prompt"),
"platform": i.get("platform"),
"insertime": i.get("insertime"),
"rank": i.get("rank")
}
)
plat_form_map = {
"DP": "deepseek网页版",
"DB": "豆包网页版",
"TXYB": "腾讯元宝",
"TYQW": "通义千问",
"KIMI": "kimi",
"WXYY": "文心一言",
"BDAI": "百度ai",
"DYAI": "抖音ai",
"DOUBA": "豆包安卓版",
"DPA": "deepseek安卓版",
}
def build_row(i):
reqId = i.get("reqId")
prompt = i.get("prompt")
platform = i.get("platform")
insertime = i.get("insertime")
rank = i.get("rank")
file = f"geo_snipaste/{pt}/{platform}/{reqId}/text.json"
context_path = f"geo_snipaste/{pt}/{platform}/{reqId}/context.txt"
quote_path = f"geo_snipaste/{pt}/{platform}/{reqId}/quote.txt"
think = f"geo_snipaste/{pt}/{platform}/{reqId}/think.txt"
png_url = f'https://tcdn.aidso.com/geo_snipaste/{pt}/{platform}/{reqId}/png.png'
text_json_str = tos_utils.get_string_from_tos(file) or "{}"
return (
prompt,
plat_form_map.get(platform),
tos_utils.get_string_from_tos(context_path),
tos_utils.get_string_from_tos(think),
tos_utils.get_string_from_tos(quote_path),
json.loads(text_json_str).get('share_url', ''),
png_url,
timestamp_to_datetime(insertime),
get_context(png_url),
rank
)
batch_data = [None] * len(task_list)
with ThreadPoolExecutor(max_workers=50) as executor:
future_map = {
executor.submit(build_row, i): index
for index, i in enumerate(task_list)
}
for future in as_completed(future_map):
index = future_map[future]
try:
batch_data[index] = future.result()
except Exception as e:
logger.exception(f"导出数据失败 index={index}, task={task_list[index]}, err={e}")
batch_data[index] = None
# 去掉失败的数据
batch_data = [row for row in batch_data if row is not None]
logger.success(f"时间{pt}--{cn}---导出数据{len(batch_data)}条")
out_path = os.path.join(EXCEL_DIR, f"美团-{pt}-{platform}-{cn}.xlsx")
excel_writer = ExcelWriter(out_path)
excel_writer.write_batch_rows(batch_data)
excel_writer.close()
return out_path
def upload_file(file_path, file_name):
url = "https://open.feishu.cn/open-apis/im/v1/files"
headers = {"Authorization": f"Bearer {get_tenant_access_token(APP_ID, APP_SECRET)}"}
files = {
"file": (
file_name,
open(file_path, "rb"),
"application/vnd.openxmlformats-officedocument.spreadsheetml.sheet"
)
}
data = {
"file_type": "xlsx",
"file_name": file_name
}
try:
resp = requests.post(url=url, headers=headers, files=files, data=data, timeout=30)
rj = resp.json()
return (rj.get("data") or {}).get("file_key")
except FileNotFoundError as e:
logger.success(f"错误:文件不存在 - {e}")
except Exception as e:
logger.success(f"请求失败:{e}")
finally:
try:
f = files.get("file")[1]
if f:
f.close()
except Exception:
pass
return None
def chat_message(file_key):
url = "https://open.feishu.cn/open-apis/im/v1/messages?receive_id_type=chat_id"
payload = json.dumps({
"content": f'{{"file_key":"{file_key}"}}',
"msg_type": "file",
"receive_id": container_id,
"uuid": str(uuid.uuid4())
})
headers = {
"Authorization": f"Bearer {get_tenant_access_token(APP_ID, APP_SECRET)}",
"Content-Type": "application/json"
}
resp = requests.request("POST", url, headers=headers, data=payload, timeout=30)
logger.success(resp.text)
def process_cn_after_tasks_done(pt: str, cn: int,platform):
redis_client = init_redis4()
logger.success(f"[pt={pt} cn={cn}] 等待 {platform} 跑空...")
wait_mt_tasks_done(redis_client,platform)
logger.success(f"[pt={pt} cn={cn}] {platform} 已为空,初步完成")
for i in range(3):
diff_result = bh_utils.query_data(
f"select a1.reqId,a1.prompt,a1.`rank` from (select reqId,prompt,`rank` from geo_feishu_snipaste where pt = {pt} and cn = {cn} and platform ='{platform}') as a1 left join geo_third_task as a2 on a1.reqId = a2.reqId where a2.reqId is null")
if not diff_result:
logger.success(f"[pt={pt} cn={cn}] {platform} 第{i + 1}轮检查已无缺失")
break
for row in diff_result:
row["thinking_enabled"] = "1"
values = [json.dumps(row, ensure_ascii=False) for row in diff_result]
redis_client.sadd(f"mt:feishu_snipaste:{platform}", *values)
logger.success(f"[pt={pt} cn={cn} {platform}] 缺失数据开始补充")
wait_mt_tasks_done(redis_client, platform)
logger.success(f"[pt={pt} cn={cn} {platform}] 缺失数据补充完成")
time.sleep(120)
logger.success(f"[pt={pt} cn={cn}] {platform} 已为空,开始汇总导出/发送")
# 1) 解析原始 json
result = bh_utils.query_data(f"select * from geo_feishu_snipaste where cn = {cn} and pt = {pt} and platform ='{platform}' ") or []
for i in result:
req_id = i.get("reqId")
if not req_id:
continue
file = f"geo_snipaste/{pt}/{platform}/{req_id}/text.json"
try:
raw = tos_utils.get_string_from_tos(file)
if not raw:
continue
content = json.loads(raw).get("content")
if content:
process_func = PLATFORM_PROCESS_MAP.get(platform)
process_func(file, content)
# doubao_process_original_data(file, content)
except Exception as e:
logger.success(f"[pt={pt} cn={cn}] 处理 {file} 失败: {e}")
# 2) 导出 Excel
out_path = to_excel2(pt, cn,platform)
# 3) 上传 + 发飞书
file_name = os.path.basename(out_path)
file_key = upload_file(out_path, file_name)
if file_key:
chat_message(file_key)
logger.success(f"[pt={pt} cn={cn} {platform}] 已发送飞书文件: {file_name}")
else:
logger.success(f"[pt={pt} cn={cn} {platform}] 上传失败,未发送飞书消息")
def send_task(prompt_list,pt,count,platform_list):
data_list = []
for platform in platform_list:
for cn in range(count):
for index,p in enumerate(prompt_list):
req_id = str(uuid.uuid4())
data_list.append({
"reqId": req_id,
"prompt": p,
"platform": platform,
"cn": cn,
"pt": pt,
"rank": index
})
logger.success(f"[init_task] 已写入任务: pt={pt}, prompts={len(prompt_list)}, rows={len(data_list)}")
bh_utils.insert_data("geo_feishu_snipaste", data_list)
def push_platform_tasks(pt_today: str, platform: str):
"""
单个平台投递任务:
cn=0/1/2 顺序写入该平台自己的 Redis 队列
"""
r = init_redis4()
redis_key = f"mt:feishu_snipaste:{platform}"
for cn in (0, 1, 2):
try:
rows = bh_utils.query_data(
f"""
SELECT reqId, prompt, rank
FROM geo_feishu_snipaste
WHERE cn = {cn}
AND pt = {pt_today}
AND platform = '{platform}'
ORDER BY rank ASC
"""
) or []
for row in rows:
row["thinking_enabled"] = "1"
values = [json.dumps(row, ensure_ascii=False) for row in rows]
if values:
r.sadd(redis_key, *values)
process_cn_after_tasks_done(pt_today, cn,platform)
logger.success(
f"[push_platform_tasks] pt={pt_today}, platform={platform}, cn={cn}, count={len(values)}, redis_key={redis_key}"
)
except Exception as e:
logger.exception(
f"[push_platform_tasks] pt={pt_today}, platform={platform}, cn={cn} 执行失败: {e}"
)
continue
def web_hook_alone(count,pt,tos_pt,platform):
redis_client = init_redis4()
redis_key = f"mt:feishu_snipaste:{platform}"
for cn in range(count):
rows = bh_utils.query_data(
f"SELECT reqId, prompt, rank FROM geo_feishu_snipaste WHERE cn = {cn} AND pt = {pt} AND platform = '{platform}'ORDER BY rank ASC"
) or []
for row in rows:
row["thinking_enabled"] = "1"
values = [json.dumps(row, ensure_ascii=False) for row in rows]
if values:
redis_client.sadd(redis_key, *values)
logger.success(f"[pt={pt} cn={cn} {platform}] 缺失数据开始补充")
wait_mt_tasks_done(redis_client, platform)
logger.success(f"[pt={pt} cn={cn} {platform}] 缺失数据补充完成")
for i in range(3):
diff_result = bh_utils.query_data(
f"select a1.reqId,a1.prompt,a1.`rank` from (select reqId,prompt,`rank` from geo_feishu_snipaste where pt = {pt} and cn = {cn} and platform ='{platform}') as a1 left join geo_third_task as a2 on a1.reqId = a2.reqId where a2.reqId is null")
for row in diff_result:
row["thinking_enabled"] = "1"
values = [json.dumps(row, ensure_ascii=False) for row in diff_result]
if not values:
logger.success(f"[pt={pt} cn={cn}] 第{i + 1}轮检查已无缺失")
break
redis_client.sadd(redis_key, *values)
logger.success(f"[pt={pt} cn={cn} {platform}] 缺失数据开始补充")
wait_mt_tasks_done(redis_client, platform)
logger.success(f"[pt={pt} cn={cn} {platform}] 缺失数据补充完成")
time.sleep(180)
logger.success(f"[pt={pt} cn={cn} {platform}] 已为空,开始汇总导出/发送")
result = bh_utils.query_data(f"select * from geo_feishu_snipaste where cn = {cn} and pt = {pt} and platform ='{platform}'") or []
for i in result:
req_id = i.get("reqId")
if not req_id:
continue
file = f"geo_snipaste/{tos_pt}/{platform}/{req_id}/text.json"
try:
raw = tos_utils.get_string_from_tos(file)
if not raw:
continue
content = json.loads(raw).get("content")
if content:
process_func = PLATFORM_PROCESS_MAP.get(platform)
process_func(file, content)
except Exception as e:
logger.success(f"[pt={pt} cn={cn}] 处理 {file} 失败: {e}")
query_list_log = bh_utils.query_data(
f"select * from geo_feishu_snipaste where cn = {cn} and pt = {pt} and platform ='{platform}'"
)
uuid_list = []
for i in query_list_log:
uuid_list.append(i.get('reqId'))
uuid_str = ",".join([f"'{uuid}'" for uuid in uuid_list])
query_list = bh_utils.query_data(
f"SELECT DISTINCT reqId,prompt, platform,insertime,rank FROM geo_third_task WHERE reqId IN ({uuid_str}) order by rank asc"
)
task_list = []
for i in query_list:
task_list.append(
{
"reqId": i.get("reqId"),
"prompt": i.get("prompt"),
"platform": i.get("platform"),
"insertime": i.get("insertime"),
"rank": i.get("rank")
}
)
plat_form_map = {
"DP": "deepseek网页版",
"DB": "豆包网页版",
"TXYB": "腾讯元宝",
"TYQW": "通义千问",
"KIMI": "kimi",
"WXYY": "文心一言",
"BDAI": "百度ai",
"DYAI": "抖音ai",
"DOUBA": "豆包安卓版",
"DPA": "deepseek安卓版",
}
def build_row(i):
reqId = i.get("reqId")
prompt = i.get("prompt")
platform = i.get("platform")
insertime = i.get("insertime")
rank = i.get("rank")
file = f"geo_snipaste/{tos_pt}/{platform}/{reqId}/text.json"
context_path = f"geo_snipaste/{tos_pt}/{platform}/{reqId}/context.txt"
quote_path = f"geo_snipaste/{tos_pt}/{platform}/{reqId}/quote.txt"
think = f"geo_snipaste/{tos_pt}/{platform}/{reqId}/think.txt"
png_url = f'https://tcdn.aidso.com/geo_snipaste/{tos_pt}/{platform}/{reqId}/png.png'
text_json_str = tos_utils.get_string_from_tos(file) or "{}"
return (
prompt,
plat_form_map.get(platform),
tos_utils.get_string_from_tos(context_path),
tos_utils.get_string_from_tos(think),
tos_utils.get_string_from_tos(quote_path),
json.loads(text_json_str).get('share_url', ''),
png_url,
timestamp_to_datetime(insertime),
get_context(png_url),
rank
)
batch_data = [None] * len(task_list)
with ThreadPoolExecutor(max_workers=50) as executor:
future_map = {
executor.submit(build_row, i): index
for index, i in enumerate(task_list)
}
for future in as_completed(future_map):
index = future_map[future]
try:
batch_data[index] = future.result()
except Exception as e:
logger.exception(f"导出数据失败 index={index}, task={task_list[index]}, err={e}")
batch_data[index] = None
# 去掉失败的数据
batch_data = [row for row in batch_data if row is not None]
out_path = os.path.join(EXCEL_DIR, f"美团-{pt}-{platform}-{cn}.xlsx")
excel_writer = ExcelWriter(out_path)
excel_writer.write_batch_rows(batch_data)
excel_writer.close()
file_name = os.path.basename(out_path)
file_key = upload_file(out_path, file_name)
if file_key:
chat_message(file_key)
logger.success(f"[pt={pt} cn={cn}] 已发送飞书文件: {file_name}")
else:
logger.success(f"[pt={pt} cn={cn}] 上传失败,未发送飞书消息")
def web_hook(promp_list,platform_list):
count = 1
pt = datetime.now().strftime("%Y%m%d%H%M")
tos_pt = datetime.now().strftime("%Y%m%d")
send_task(promp_list, pt, count,platform_list)
with ThreadPoolExecutor(max_workers=len(platform_list)) as executor:
future_map = {
executor.submit(web_hook_alone,count,pt,tos_pt,platform): platform
for platform in platform_list
}
for future in as_completed(future_map):
platform = future_map[future]
try:
future.result()
logger.success(f"[web_hook] 平台处理完成: pt={pt}, platform={platform}")
except Exception as e:
logger.exception(f"[web_hook] 平台处理失败: pt={pt}, platform={platform}, err={e}")
def run_once_for_pt(pt_today: str):
"""
四个平台并发投递:
DB / DP / TXYB / TYQW 同时跑
每个平台内部 cn=0/1/2 顺序跑
"""
platform_list = ["DB", "DP", "TXYB", "TYQW"]
with ThreadPoolExecutor(max_workers=4) as executor:
future_map = {
executor.submit(push_platform_tasks, pt_today, platform): platform
for platform in platform_list
}
for future in as_completed(future_map):
platform = future_map[future]
try:
future.result()
logger.success(f"[run_once_for_pt] 平台投递完成: pt={pt_today}, platform={platform}")
except Exception as e:
logger.exception(f"[run_once_for_pt] 平台投递失败: pt={pt_today}, platform={platform}, err={e}")
def init_tomorrow_tasks():
"""
每天晚上 23:00 执行:
初始化明天各个平台的任务
"""
tomorrow = datetime.now().date() + timedelta(days=1)
pt_tomorrow = tomorrow.strftime("%Y%m%d")
platform_list = ["DB", "DP", "TXYB", "TYQW"]
for p in platform_list:
try:
init_task_once(pt_tomorrow, p)
logger.success(f"[init_tomorrow_tasks] 初始化成功 pt={pt_tomorrow}, platform={p}")
except Exception as e:
logger.exception(f"[init_tomorrow_tasks] 初始化失败 pt={pt_tomorrow}, platform={p}, err={e}")
def run_task():
pt_today = datetime.now().date().strftime("%Y%m%d")
run_once_for_pt(pt_today)
if __name__ == "__main__":
scheduler = BlockingScheduler(timezone="Asia/Shanghai")
# 每天晚上 23:00 执行
scheduler.add_job(
init_tomorrow_tasks,
trigger="cron",
hour=23,
minute=0,
second=0,
id="init_tomorrow_tasks",
replace_existing=True,
max_instances=1,
coalesce=True,
)
# 每天晚上 00:10 执行
scheduler.add_job(
run_task,
trigger="cron",
hour=0,
minute=10,
second=0,
id="run_task",
replace_existing=True,
max_instances=1,
coalesce=True,
)
logger.success(
"定时任务注册完成:"
"init_tomorrow_tasks(每天23:00), "
"run_task(每天0:10),"
)
scheduler.start()
# platform_list = ["DB", "DP", "TXYB", "TYQW"]
# for platform in platform_list:
# init_task('20260610', platform)
# run_once_for_pt('20260610')
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