304 lines
17 KiB
Python
304 lines
17 KiB
Python
import asyncio
|
|
from datetime import datetime, timezone, timedelta
|
|
import json
|
|
|
|
from sqlalchemy import select, delete
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from app.models.base import async_session
|
|
from app.models.material_cost import MaterialCost
|
|
from app.models.user_oauth import UserOAuth
|
|
from app.models.user_oauth_account import UserOAuthAccount
|
|
from app.models.resources_material import ResourcesMaterial
|
|
from app.utils.id_gen import generate_id
|
|
from app.utils.douyinApi import DouyinApi
|
|
from app.utils.logger import get_logger
|
|
|
|
logger = get_logger("material_consumption_task", "material_consumption")
|
|
|
|
douyin_api = DouyinApi()
|
|
|
|
|
|
def _convert_numeric(value):
|
|
"""Convert string numeric values to appropriate Python types."""
|
|
if value is None or value == "":
|
|
return None
|
|
try:
|
|
# Try int first
|
|
return int(value)
|
|
except ValueError:
|
|
try:
|
|
# Then try float
|
|
return float(value)
|
|
except ValueError:
|
|
# Return original value if conversion fails
|
|
return None
|
|
|
|
|
|
class MaterialConsumptionQueue:
|
|
def __init__(self):
|
|
self.queue: asyncio.Queue[dict] = asyncio.Queue()
|
|
self.running = False
|
|
|
|
async def enqueue(self, task_data: dict):
|
|
"""Add a task to the queue."""
|
|
await self.queue.put(task_data)
|
|
|
|
async def run(self):
|
|
"""Main processing loop."""
|
|
self.running = True
|
|
logger.info("Material consumption queue started")
|
|
|
|
while self.running:
|
|
try:
|
|
task_data = await asyncio.wait_for(self.queue.get(), timeout=5.0)
|
|
except asyncio.TimeoutError:
|
|
continue
|
|
|
|
try:
|
|
await self._process(task_data)
|
|
except Exception as e:
|
|
logger.error(f"Error processing material consumption task: {e}")
|
|
finally:
|
|
self.queue.task_done()
|
|
|
|
logger.info("Material consumption queue stopped")
|
|
|
|
async def _process(self, task_data: dict):
|
|
"""Process a single consumption update task."""
|
|
oauth_id = task_data.get("oauth_id")
|
|
advertiser_id = task_data.get("advertiser_id")
|
|
date = task_data.get("date")
|
|
|
|
logger.info(f"Processing consumption update for oauth_id={oauth_id}, advertiser_id={advertiser_id}, date={date}")
|
|
|
|
try:
|
|
result = await _fetch_and_save_consumption(oauth_id, advertiser_id, date)
|
|
if result.get("success"):
|
|
logger.info(f"Successfully updated consumption for advertiser {advertiser_id} on {date}")
|
|
else:
|
|
logger.error(f"Failed to update consumption for advertiser {advertiser_id} on {date}: {result.get('error')}")
|
|
except Exception as e:
|
|
logger.error(f"Exception processing consumption for advertiser {advertiser_id} on {date}: {e}")
|
|
|
|
def stop(self):
|
|
"""Stop the queue."""
|
|
self.running = False
|
|
|
|
|
|
async def _fetch_and_save_consumption(oauth_id: str, advertiser_id: str, date: str) -> dict:
|
|
async with async_session() as db:
|
|
try:
|
|
consume_date_obj = datetime.strptime(date, "%Y-%m-%d").date()
|
|
|
|
# 查询当前授权下系统上传的素材ID列表
|
|
material_result = await db.execute(
|
|
select(ResourcesMaterial.upload_id)
|
|
.where(
|
|
ResourcesMaterial.oauth_id == oauth_id,
|
|
ResourcesMaterial.advertiser_id == advertiser_id,
|
|
ResourcesMaterial.deleted_at.is_(None),
|
|
ResourcesMaterial.upload_id.isnot(None),
|
|
ResourcesMaterial.upload_id != "",
|
|
)
|
|
)
|
|
system_material_ids = [str(row[0]) for row in material_result.all()]
|
|
|
|
if not system_material_ids:
|
|
return {
|
|
"success": True,
|
|
"message": "没有系统上传的素材"
|
|
}
|
|
|
|
#强制删除旧数据,真实删除,不要软删除
|
|
await db.execute(
|
|
delete(MaterialCost).where(
|
|
MaterialCost.oauth_id == oauth_id,
|
|
MaterialCost.advertiser_id == advertiser_id,
|
|
MaterialCost.consume_date == consume_date_obj,
|
|
MaterialCost.deleted_at.is_(None),
|
|
)
|
|
)
|
|
|
|
BATCH_SIZE = 50
|
|
total_saved = 0
|
|
total_requests = 0
|
|
|
|
# 分批请求,每次请求50个素材
|
|
for i in range(0, len(system_material_ids), BATCH_SIZE):
|
|
batch_materials = system_material_ids[i:i+BATCH_SIZE]
|
|
total_requests += 1
|
|
|
|
params = {
|
|
"dimensions": json.dumps(["ad_platform_material_name","image_mode","material_id","stat_time_day"]),
|
|
"advertiser_id": int(advertiser_id),
|
|
"metrics": json.dumps([
|
|
"stat_cost","show_cnt","cpm_platform","click_cnt","ctr","cpc_platform","convert_cnt","conversion_cost","conversion_rate","deep_convert_cnt","deep_convert_cost","deep_convert_rate","click_start_cnt","click_start_cost","download_finish_rate",
|
|
"install_finish_cnt","install_finish_cost","install_finish_rate","active","active_cost",
|
|
"active_rate","active_register","active_register_cost","active_register_rate","game_addiction",
|
|
"game_addiction_cost","game_addiction_rate","attribution_next_day_open_cnt","attribution_next_day_open_cost","attribution_next_day_open_rate","next_day_open","active_pay","active_pay_cost","active_pay_rate",
|
|
"game_pay_count","game_pay_cost","in_app_uv","in_app_detail_uv","in_app_cart","in_app_pay",
|
|
"in_app_order","attribution_billing_game_in_app_ltv_1day","attribution_billing_game_in_app_roi_1day","phone","form",
|
|
"form_submit","map","button","view","download_start","qq","vote","lottery","message",
|
|
"redirect","shopping","consult","consult_effective","phone_confirm","phone_connect","phone_effective",
|
|
"redirect_to_shop","coupon_single_page","poi_address_click","poi_collect","customer_effective",
|
|
"attribution_customer_effective","attribution_customer_effective_cost","attribution_clue_pay_succeed",
|
|
"attribution_clue_pay_succeed_cost","attribution_clue_interflow","attribution_clue_interflow_cost","attribution_clue_high_intention","attribution_clue_high_intention_cost",
|
|
"consult_clue","clue_message_count","attribution_work_wechat_unfriend_count","attribution_clue_connected_count","attribution_clue_connected_cost","attribution_clue_connected_rate","attribution_form","form_and_submit_count","form_and_submit_cost","intention_form_and_submit_count","intention_form_and_submit_cost","attribution_micro_game_0d_ltv","attribution_micro_game_3d_ltv","attribution_micro_game_7d_ltv","attribution_micro_game_0d_roi","attribution_micro_game_3d_roi","attribution_micro_game_7d_roi","attribution_game_in_app_ltv_1day","attribution_game_in_app_roi_1day","active_pay_intra_day_count","active_pay_intra_day_cost","active_pay_intra_day_rate","first_pay_intra_24hour_amount","loan_completion","loan_completion_cost","loan_completion_rate","pre_loan_credit","pre_loan_credit_cost","loan_credit_cost","loan_credit","loan_credit_rate","in_wechat_pay_count","unfollow_in_wechat_count","loan","loan_cost","loan_rate","open_account_count","withdraw_m2_count","in_app_order_gmv","in_app_order_roi","in_app_pay_gmv","in_app_pay_roi","first_rental_order_count","commute_first_pay_count","first_order_count","submit_certification_count","approval_count","total_play","play_duration_3s","valid_play","valid_play_cost","valid_play_rate","valid_play_of_mille","valid_play_cost_of_mille","average_play_time_per_play","play_over_rate","dy_like","dy_comment","dy_share","report_cnt","location_click","dislike_cnt","dy_home_visited","dy_follow","message_action","click_landing_page","click_shopwindow","click_website","click_call_dy","click_download","luban_live_enter_cnt","luban_live_follow_cnt","luban_live_share_cnt","luban_live_comment_cnt","luban_live_gift_cnt","luban_live_gift_amount","click_call_cnt","click_counsel","a3_ask_count","a3_ask_cost"
|
|
]),
|
|
"filters": json.dumps([
|
|
{
|
|
"field": "stat_cost",
|
|
"type": 3,
|
|
"operator": 4,
|
|
"values": ["0"]
|
|
},
|
|
{
|
|
"field": "material_id",
|
|
"type": 2,
|
|
"operator": 7,
|
|
"values": batch_materials
|
|
}
|
|
]),
|
|
"start_time": date,
|
|
"end_time": date,
|
|
"order_by": json.dumps([
|
|
{
|
|
"field": "stat_cost",
|
|
"type": "ASC"
|
|
}
|
|
])
|
|
}
|
|
|
|
response = await douyin_api.get_material_cost(oauth_id, params, request_count=3)
|
|
|
|
if response.get("code") != 0:
|
|
logger.warning(f"获取素材消耗失败(批次 {i//BATCH_SIZE + 1}): {response.get('message', '未知错误')}")
|
|
continue
|
|
|
|
data = response.get("data", {})
|
|
list_data = data.get("rows", [])
|
|
|
|
for item in list_data:
|
|
dimensions = item.get("dimensions", {})
|
|
metrics = item.get("metrics", {})
|
|
material_id = dimensions.get("material_id", "")
|
|
|
|
cost = MaterialCost(
|
|
id=generate_id(),
|
|
oauth_id=oauth_id,
|
|
advertiser_id=str(advertiser_id),
|
|
material_id=material_id,
|
|
consume_date=consume_date_obj,
|
|
stat_cost=_convert_numeric(metrics.get("stat_cost")),
|
|
show_cnt=_convert_numeric(metrics.get("show_cnt")),
|
|
cpm_platform=_convert_numeric(metrics.get("cpm_platform")),
|
|
click_cnt=_convert_numeric(metrics.get("click_cnt")),
|
|
ctr=_convert_numeric(metrics.get("ctr")),
|
|
cpc_platform=_convert_numeric(metrics.get("cpc_platform")),
|
|
convert_cnt=_convert_numeric(metrics.get("convert_cnt")),
|
|
conversion_cost=_convert_numeric(metrics.get("conversion_cost")),
|
|
conversion_rate=_convert_numeric(metrics.get("conversion_rate")),
|
|
deep_convert_cnt=_convert_numeric(metrics.get("deep_convert_cnt")),
|
|
deep_convert_cost=_convert_numeric(metrics.get("deep_convert_cost")),
|
|
deep_convert_rate=_convert_numeric(metrics.get("deep_convert_rate")),
|
|
active=_convert_numeric(metrics.get("active")),
|
|
active_cost=_convert_numeric(metrics.get("active_cost")),
|
|
active_rate=_convert_numeric(metrics.get("active_rate")),
|
|
active_register=_convert_numeric(metrics.get("active_register")),
|
|
active_register_cost=_convert_numeric(metrics.get("active_register_cost")),
|
|
active_register_rate=_convert_numeric(metrics.get("active_register_rate")),
|
|
attribution_next_day_open_cnt=_convert_numeric(metrics.get("attribution_next_day_open_cnt")),
|
|
attribution_next_day_open_cost=_convert_numeric(metrics.get("attribution_next_day_open_cost")),
|
|
attribution_next_day_open_rate=_convert_numeric(metrics.get("attribution_next_day_open_rate")),
|
|
active_pay=_convert_numeric(metrics.get("active_pay")),
|
|
active_pay_cost=_convert_numeric(metrics.get("active_pay_cost")),
|
|
active_pay_rate=_convert_numeric(metrics.get("active_pay_rate")),
|
|
phone=_convert_numeric(metrics.get("phone")),
|
|
form=_convert_numeric(metrics.get("form")),
|
|
download_start=_convert_numeric(metrics.get("download_start")),
|
|
form_submit=_convert_numeric(metrics.get("form_submit")),
|
|
button=_convert_numeric(metrics.get("button")),
|
|
view=_convert_numeric(metrics.get("view")),
|
|
message=_convert_numeric(metrics.get("message")),
|
|
consult=_convert_numeric(metrics.get("consult")),
|
|
consult_effective=_convert_numeric(metrics.get("consult_effective")),
|
|
shopping=_convert_numeric(metrics.get("shopping")),
|
|
customer_effective=_convert_numeric(metrics.get("customer_effective")),
|
|
attribution_game_in_app_ltv_1day=_convert_numeric(metrics.get("attribution_game_in_app_ltv_1day")),
|
|
attribution_game_in_app_roi_1day=_convert_numeric(metrics.get("attribution_game_in_app_roi_1day")),
|
|
loan_completion=_convert_numeric(metrics.get("loan_completion")),
|
|
loan_completion_cost=_convert_numeric(metrics.get("loan_completion_cost")),
|
|
loan_completion_rate=_convert_numeric(metrics.get("loan_completion_rate")),
|
|
loan_credit=_convert_numeric(metrics.get("loan_credit")),
|
|
loan_credit_cost=_convert_numeric(metrics.get("loan_credit_cost")),
|
|
loan_credit_rate=_convert_numeric(metrics.get("loan_credit_rate")),
|
|
in_app_order_gmv=_convert_numeric(metrics.get("in_app_order_gmv")),
|
|
in_app_order_roi=_convert_numeric(metrics.get("in_app_order_roi")),
|
|
in_app_pay_gmv=_convert_numeric(metrics.get("in_app_pay_gmv")),
|
|
in_app_pay_roi=_convert_numeric(metrics.get("in_app_pay_roi")),
|
|
total_play=_convert_numeric(metrics.get("total_play")),
|
|
valid_play=_convert_numeric(metrics.get("valid_play")),
|
|
valid_play_cost=_convert_numeric(metrics.get("valid_play_cost")),
|
|
valid_play_rate=_convert_numeric(metrics.get("valid_play_rate")),
|
|
valid_play_of_mille=_convert_numeric(metrics.get("valid_play_of_mille")),
|
|
valid_play_cost_of_mille=_convert_numeric(metrics.get("valid_play_cost_of_mille")),
|
|
average_play_time_per_play=_convert_numeric(metrics.get("average_play_time_per_play")),
|
|
play_over_rate=_convert_numeric(metrics.get("play_over_rate")),
|
|
dy_like=_convert_numeric(metrics.get("dy_like")),
|
|
dy_comment=_convert_numeric(metrics.get("dy_comment")),
|
|
dy_share=_convert_numeric(metrics.get("dy_share")),
|
|
report_cnt=_convert_numeric(metrics.get("report_cnt")),
|
|
)
|
|
db.add(cost)
|
|
total_saved += 1
|
|
|
|
await db.commit()
|
|
|
|
return {
|
|
"success": True,
|
|
"message": f"成功更新 {total_saved} 条记录,共请求 {total_requests} 批次,系统素材总数 {len(system_material_ids)}",
|
|
"count": total_saved
|
|
}
|
|
|
|
except Exception as e:
|
|
await db.rollback()
|
|
return {
|
|
"success": False,
|
|
"error": str(e)
|
|
}
|
|
|
|
#异步获取所有要更新的广告主,素材消耗
|
|
async def sync_all_advertisers_consumption(date: str = None):
|
|
"""Sync consumption data for all advertisers."""
|
|
if date is None:
|
|
date = (datetime.now(timezone.utc) - timedelta(days=1)).strftime("%Y-%m-%d")
|
|
|
|
async with async_session() as db:
|
|
result = await db.execute(
|
|
select(UserOAuth.id, UserOAuthAccount.advertiser_id)
|
|
.join(UserOAuthAccount, UserOAuth.id == UserOAuthAccount.oauth_id)
|
|
.join(ResourcesMaterial, UserOAuthAccount.oauth_id == ResourcesMaterial.oauth_id)
|
|
.where(
|
|
UserOAuth.deleted_at.is_(None),
|
|
UserOAuthAccount.deleted_at.is_(None),
|
|
ResourcesMaterial.deleted_at.is_(None),
|
|
ResourcesMaterial.material_id.isnot(None),
|
|
ResourcesMaterial.material_id != "",
|
|
)
|
|
.distinct()
|
|
)
|
|
oauth_advertiser_pairs = result.all()
|
|
|
|
for oauth_id, advertiser_id in oauth_advertiser_pairs:
|
|
await material_consumption_queue.enqueue({
|
|
"oauth_id": oauth_id,
|
|
"advertiser_id": advertiser_id,
|
|
"date": date,
|
|
})
|
|
logger.info(f"Enqueued consumption sync for oauth_id={oauth_id}, advertiser_id={advertiser_id}, date={date}")
|
|
|
|
return {"message": f"已将 {len(oauth_advertiser_pairs)} 个广告主的消耗更新任务加入队列"}
|
|
|
|
|
|
material_consumption_queue = MaterialConsumptionQueue() |