Commit ea577486 authored by Yaowentong's avatar Yaowentong

success 修复

parent 1efc8c84
......@@ -14,13 +14,13 @@ redis_client = init_redis()
QUEUE_KEY = "geo:task_commit:list"
# 每轮最多拉取任务数量
BATCH_SIZE = 1000
BATCH_SIZE = 100
# 每 60 秒执行一次
INTERVAL_SECONDS = 60
# 每 30 秒执行一次
INTERVAL_SECONDS = 30
# 每轮内部并发数
CONCURRENT_WORKERS = 40
CONCURRENT_WORKERS = 20
def parse_task(raw):
......@@ -128,8 +128,8 @@ if __name__ == "__main__":
seconds=INTERVAL_SECONDS,
id="consume_geo_task_commit_queue",
max_instances=5, # 允许最多 5 个调度批次同时跑
coalesce=False, # 不合并错过的执行
misfire_grace_time=60
coalesce=True,
misfire_grace_time=300
)
logger.success(
......
......@@ -1029,6 +1029,11 @@ def get_platforms_q(brand_library_id, user_id, begin_time, end_time,question_lis
# 使用示例
# =========================
if __name__ == "__main__":
result = get_req_id(18156037075,'2026-06-24','2026-06-24')
req_list = []
for i in result:
req_list.append(i.get('req_id'))
print(req_list)
# print(cha_report_user(3738753854605260))
# keyword_list = ['618']
# start = '20250518'
......@@ -1048,19 +1053,19 @@ if __name__ == "__main__":
# zhou_report(phone, begin, end, b)
# # qian_report(phone,begin,end,b)
# qian_report(phone,begin,end,b)
brand_library_id = 2063180776253702144
user_id = 2063177785609949184
begin = '2026-06-08'
end = '2026-06-08'
qu_list = ['东鹏瓷砖怎么样','东鹏瓷砖质量好不好','东鹏瓷砖值得买吗','东鹏瓷砖是几线品牌','东鹏瓷砖和马可波罗哪个好','东鹏和冠珠瓷砖哪个好','东鹏和蒙娜丽莎哪个好','东鹏控股是做什么的','东鹏控股和东鹏饮料什么关系','瓷砖十大品牌有哪些','2026年瓷砖品牌排行榜前十名','中国瓷砖品牌排名前十','瓷砖一线品牌排名','国内瓷砖品牌排行榜','高端瓷砖品牌排行榜','高端瓷砖有哪些品牌','瓷砖头部品牌有哪些','大平层用什么瓷砖品牌','设计师推荐的高端瓷砖品牌','5A瓷砖品牌推荐哪个好','5A瓷砖品牌排行','5A国标瓷砖什么品牌好','5A瓷砖是什么标准','5A认证瓷砖推荐','品质好的瓷砖品牌有哪些','什么品牌的瓷砖品质最好','选好瓷砖认准什么品牌','瓷砖哪个牌子好','瓷砖品牌推荐','装修选什么瓷砖品牌好','瓷砖什么牌子质量好','2026年瓷砖品牌排行榜','瓷砖怎么选不踩坑','买瓷砖主要看哪几个指标','好瓷砖的标准是什么','瓷砖选购避坑指南','金丝绒瓷砖哪个牌子好','金丝绒瓷砖值得买吗','金丝绒瓷砖怎么样','木纹砖哪个牌子好','木纹砖推荐哪个品牌','木纹砖怎么选品牌','木纹砖和木地板哪个好','木纹砖品牌排行','莱姆石瓷砖哪个牌子好','莱姆石瓷砖品牌推荐','莱姆石瓷砖怎么选','什么品牌的莱姆石瓷砖好','客厅瓷砖什么牌子好','客厅铺什么瓷砖好看又耐用','客厅瓷砖怎么选品牌','客厅瓷砖品牌推荐','客厅用什么瓷砖显高级','厨房瓷砖什么牌子好','厨房用什么瓷砖好打理','厨房防油污瓷砖推荐哪个牌子','厨房瓷砖品牌推荐2026','厨房抗菌瓷砖推荐哪个品牌','卫生间瓷砖什么牌子好','卫生间瓷砖推荐哪个品牌','浴室防滑瓷砖哪个品牌好','卫生间防滑抗菌瓷砖推荐品牌','浴室抗菌瓷砖什么品牌好','全屋通铺瓷砖什么品牌好','全屋通铺瓷砖推荐哪个牌子','全屋瓷砖用什么品牌好','家里全屋铺瓷砖选什么牌子','全屋通铺瓷砖品牌排行','防滑瓷砖哪个牌子好','防滑瓷砖品牌推荐','好打理的瓷砖推荐哪个品牌','防污瓷砖什么品牌好','耐脏好清洁的瓷砖品牌推荐','耐磨瓷砖什么品牌好','不容易刮花的瓷砖推荐什么品牌','中古风装修用什么瓷砖品牌好','法式风格瓷砖推荐什么品牌','奶油风瓷砖什么品牌好','现代简约风格瓷砖推荐什么品牌','新中式瓷砖用什么品牌好','原木风瓷砖推荐哪个品牌','高级感装修瓷砖用什么品牌','岩板什么品牌好','岩板品牌推荐','750x1500地砖什么品牌好','岩板品牌排行','柔光砖什么品牌好','哑光瓷砖推荐哪个品牌','哑光砖品牌排名','哑光瓷砖什么品牌好','柔光砖品牌排行','哑光砖推荐哪个品牌耐脏','国家级建筑用的瓷砖是什么牌子','大型工程项目用什么瓷砖品牌','瓷砖行业有哪些上市公司','A股瓷砖上市公司有哪些','绿色建材瓷砖品牌有哪些','获得国家级绿色工厂认证的瓷砖企业有哪些','建材行业ESG表现好的企业有哪些','双碳目标下有哪些绿色建材瓷砖品牌']
result = []
for q in qu_list:
result.append((q,get_platforms_q(brand_library_id,user_id,begin,end,[q]).get('geo_refer_rank_overview_vos')))
# brand_name = 'AIDSO爱搜'
all_file = f"/Users/yaowentong/Desktop/dongpeng.txt"
with open(all_file, "w", encoding="utf-8") as f:
for item in result:
f.write(json.dumps(item, ensure_ascii=False) + "\n\n\n")
# brand_library_id = 2063180776253702144
# user_id = 2063177785609949184
# begin = '2026-06-08'
# end = '2026-06-08'
# qu_list = ['东鹏瓷砖怎么样','东鹏瓷砖质量好不好','东鹏瓷砖值得买吗','东鹏瓷砖是几线品牌','东鹏瓷砖和马可波罗哪个好','东鹏和冠珠瓷砖哪个好','东鹏和蒙娜丽莎哪个好','东鹏控股是做什么的','东鹏控股和东鹏饮料什么关系','瓷砖十大品牌有哪些','2026年瓷砖品牌排行榜前十名','中国瓷砖品牌排名前十','瓷砖一线品牌排名','国内瓷砖品牌排行榜','高端瓷砖品牌排行榜','高端瓷砖有哪些品牌','瓷砖头部品牌有哪些','大平层用什么瓷砖品牌','设计师推荐的高端瓷砖品牌','5A瓷砖品牌推荐哪个好','5A瓷砖品牌排行','5A国标瓷砖什么品牌好','5A瓷砖是什么标准','5A认证瓷砖推荐','品质好的瓷砖品牌有哪些','什么品牌的瓷砖品质最好','选好瓷砖认准什么品牌','瓷砖哪个牌子好','瓷砖品牌推荐','装修选什么瓷砖品牌好','瓷砖什么牌子质量好','2026年瓷砖品牌排行榜','瓷砖怎么选不踩坑','买瓷砖主要看哪几个指标','好瓷砖的标准是什么','瓷砖选购避坑指南','金丝绒瓷砖哪个牌子好','金丝绒瓷砖值得买吗','金丝绒瓷砖怎么样','木纹砖哪个牌子好','木纹砖推荐哪个品牌','木纹砖怎么选品牌','木纹砖和木地板哪个好','木纹砖品牌排行','莱姆石瓷砖哪个牌子好','莱姆石瓷砖品牌推荐','莱姆石瓷砖怎么选','什么品牌的莱姆石瓷砖好','客厅瓷砖什么牌子好','客厅铺什么瓷砖好看又耐用','客厅瓷砖怎么选品牌','客厅瓷砖品牌推荐','客厅用什么瓷砖显高级','厨房瓷砖什么牌子好','厨房用什么瓷砖好打理','厨房防油污瓷砖推荐哪个牌子','厨房瓷砖品牌推荐2026','厨房抗菌瓷砖推荐哪个品牌','卫生间瓷砖什么牌子好','卫生间瓷砖推荐哪个品牌','浴室防滑瓷砖哪个品牌好','卫生间防滑抗菌瓷砖推荐品牌','浴室抗菌瓷砖什么品牌好','全屋通铺瓷砖什么品牌好','全屋通铺瓷砖推荐哪个牌子','全屋瓷砖用什么品牌好','家里全屋铺瓷砖选什么牌子','全屋通铺瓷砖品牌排行','防滑瓷砖哪个牌子好','防滑瓷砖品牌推荐','好打理的瓷砖推荐哪个品牌','防污瓷砖什么品牌好','耐脏好清洁的瓷砖品牌推荐','耐磨瓷砖什么品牌好','不容易刮花的瓷砖推荐什么品牌','中古风装修用什么瓷砖品牌好','法式风格瓷砖推荐什么品牌','奶油风瓷砖什么品牌好','现代简约风格瓷砖推荐什么品牌','新中式瓷砖用什么品牌好','原木风瓷砖推荐哪个品牌','高级感装修瓷砖用什么品牌','岩板什么品牌好','岩板品牌推荐','750x1500地砖什么品牌好','岩板品牌排行','柔光砖什么品牌好','哑光瓷砖推荐哪个品牌','哑光砖品牌排名','哑光瓷砖什么品牌好','柔光砖品牌排行','哑光砖推荐哪个品牌耐脏','国家级建筑用的瓷砖是什么牌子','大型工程项目用什么瓷砖品牌','瓷砖行业有哪些上市公司','A股瓷砖上市公司有哪些','绿色建材瓷砖品牌有哪些','获得国家级绿色工厂认证的瓷砖企业有哪些','建材行业ESG表现好的企业有哪些','双碳目标下有哪些绿色建材瓷砖品牌']
# result = []
# for q in qu_list:
# result.append((q,get_platforms_q(brand_library_id,user_id,begin,end,[q]).get('geo_refer_rank_overview_vos')))
# # brand_name = 'AIDSO爱搜'
# all_file = f"/Users/yaowentong/Desktop/dongpeng.txt"
# with open(all_file, "w", encoding="utf-8") as f:
# for item in result:
# f.write(json.dumps(item, ensure_ascii=False) + "\n\n\n")
# platform = ['DB']
# dao_report(phone, begin, end, brand_name, platform)
# qian_report(start, end, keyword_list, file_name)
......@@ -155,37 +155,60 @@ def task_commit():
"status": f"缺少必要字段:{', '.join(missing_fields)}",
"reqId": req_id
})
if type == 'stream':
insert_ok = commit_task(data, 'ING')
ok = submit_background_task(main_process, data)
if insert_ok and ok:
return jsonify({
"code": 200,
"msg": 'success',
"reqId": req_id
})
else:
logger.success(f"{data['reqId']}--{platform}--{prompt}--任务提交--{type}")
ret = redis_client.lpush("geo:task_commit:list",json.dumps(data, ensure_ascii=False))
if ret and ret > 0:
resp_cache_key = f"geo:task_check:resp:{req_id}"
resp = {
"code": 200,
"msg": "success",
"data": {
"status": 'ING',
"result": {}
}
# if type == 'stream':
# insert_ok = commit_task(data, 'ING')
# ok = submit_background_task(main_process, data)
# if insert_ok and ok:
# return jsonify({
# "code": 200,
# "msg": 'success',
# "reqId": req_id
# })
# else:
# logger.success(f"{data['reqId']}--{platform}--{prompt}--任务提交--{type}")
# ret = redis_client.lpush("geo:task_commit:list",json.dumps(data, ensure_ascii=False))
# if ret and ret > 0:
# resp_cache_key = f"geo:task_check:resp:{req_id}"
# resp = {
# "code": 200,
# "msg": "success",
# "data": {
# "status": 'ING',
# "result": {}
# }
# }
# if type == 'stream_batch':
# _cache_set_json(resp_cache_key, resp, 600)
# else:
# _cache_set_json(resp_cache_key, resp, 18000)
# return jsonify({
# "code": 200,
# "msg": "任务已提交",
# "reqId": req_id
# })
logger.success(f"{data['reqId']}--{platform}--{prompt}--任务提交--{type}")
ret = redis_client.lpush("geo:task_commit:list",json.dumps(data, ensure_ascii=False))
if ret and ret > 0:
resp_cache_key = f"geo:task_check:resp:{req_id}"
resp = {
"code": 200,
"msg": "success",
"data": {
"status": 'ING',
"result": {}
}
if type == 'stream_batch':
_cache_set_json(resp_cache_key, resp, 600)
else:
_cache_set_json(resp_cache_key, resp, 18000)
return jsonify({
"code": 200,
"msg": "任务已提交",
"reqId": req_id
})
}
if type == 'stream':
_cache_set_json(resp_cache_key, resp, 600)
elif type == 'stream_batch':
_cache_set_json(resp_cache_key, resp, 600)
elif type == 'batch':
_cache_set_json(resp_cache_key, resp, 18000)
return jsonify({
"code": 200,
"msg": "任务已提交",
"reqId": req_id
})
return jsonify({
"code": 503,
......@@ -211,21 +234,6 @@ def data_call_back():
req_id = task_data.get('reqId')
platform = task_data.get('platform')
logger.success(f"{req_id}--{platform}-------CALL_BACK")
# ret = redis_client8.lpush("geo:call_back:list",json.dumps(data, ensure_ascii=False))
#
# if ret and ret > 0:
# return jsonify({
# "code": 200,
# "msg": "任务已提交",
# "reqId": req_id
# })
# else:
# return jsonify({
# "code": 400,
# "msg": "提交失败",
# "reqId": req_id
# })
#
ok = submit_background_task(process_call_back, task_data, result)
if not ok:
return jsonify({
......@@ -328,76 +336,3 @@ def check_quto():
"reqId": req_id
})
@line_app.route('/api/geo/get_mention_count', methods=['POST'])
def get_mention_count():
try:
data = request.get_json(force=True) or {}
req_ids = data.get('reqIds', [])
in_sql = ",".join([f"'{i}'" for i in list(dict.fromkeys(req_ids))])
result = bh_utils.query_data(f"""
WITH t AS (
SELECT
req_id,
arrayFilter(x -> x != '',
arrayMap(x -> replaceRegexpAll(x, '(^\\s+)|(\\s+$)', ''),
splitByChar(',', ifNull(positive_mentions, '')))
) AS pos_arr,
arrayFilter(x -> x != '',
arrayMap(x -> replaceRegexpAll(x, '(^\\s+)|(\\s+$)', ''),
splitByChar(',', ifNull(negative_mentions, '')))
) AS neg_arr
FROM geo_brand_mention_list
WHERE req_id IN ({in_sql}
)
),
pos AS (
SELECT
req_id,
uniqExact(word) AS positive_distinct,
count() AS positive_total
FROM (
SELECT req_id, arrayJoin(pos_arr) AS word
FROM t
)
GROUP BY req_id
),
neg AS (
SELECT
req_id,
uniqExact(word) AS negative_distinct,
count() AS negative_total
FROM (
SELECT req_id, arrayJoin(neg_arr) AS word
FROM t
)
GROUP BY req_id
)
SELECT
a.req_id,
ifNull(p.positive_distinct, 0) AS positive_word_count,
ifNull(p.positive_total, 0) AS positive_mention_count,
ifNull(n.negative_distinct, 0) AS negative_word_count,
ifNull(n.negative_total, 0) AS negative_mention_count
FROM
(SELECT DISTINCT req_id FROM t) a
LEFT JOIN pos p ON a.req_id = p.req_id
LEFT JOIN neg n ON a.req_id = n.req_id
ORDER BY a.req_id;
""")
return jsonify({
"code": 200,
"msg": "success",
"data": result or []
})
except Exception as e:
return jsonify({
"code": 400,
"msg": "err",
"data": []
})
This diff is collapsed.
......@@ -31,7 +31,6 @@ def yuanbao_android_process_original_data(data):
# 提取并解析JSON数据
data_str = i.split("data: ")[1]
json_data = json.loads(data_str)
except (IndexError, json.JSONDecodeError):
continue # 跳过格式错误的数据
if json_data.get('type') == 'searchGuid':
......@@ -106,7 +105,7 @@ if __name__ == '__main__':
# yuanbao_android_process_original_data(file_path2)
data_list = bh_utils.query_data(f"select * from geo_commit_task where taskId = '6de35e20-4a67-4761-8a16-b91dd1f2fb7f' and platform = 'TXYBA'")
data_list = bh_utils.query_data(f"select * from geo_commit_task where taskId = 'f9018d5c-2d0a-4e6a-9476-4b92334a0715' and platform = 'TXYBA'")
# # #
# # #
......
This diff is collapsed.
......@@ -344,7 +344,6 @@ def task_queue_backlog():
if __name__ == '__main__':
logger.info("监控调度器启动")
scheduler = BlockingScheduler(timezone="Asia/Shanghai")
scheduler.add_job(
fail_task_send_feishu,
trigger='cron',
......
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