import asyncio import logging from datetime import datetime, timezone, timedelta 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 import os import json LOG_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(__file__))), "logs") os.makedirs(LOG_DIR, exist_ok=True) logger = logging.getLogger("material_consumption_task") logger.setLevel(logging.INFO) class DailyRotatingFileHandler(logging.FileHandler): def __init__(self, directory, encoding=None): self.directory = directory filename = self._get_log_filename() super().__init__(filename, encoding=encoding) def _get_log_filename(self): return os.path.join(self.directory, f"material_consumption-{datetime.now(timezone.utc).strftime('%Y-%m-%d')}.log") def emit(self, record): current_filename = self._get_log_filename() if self.baseFilename != current_filename: self.close() self.baseFilename = current_filename self.stream = self._open() super().emit(record) if not logger.handlers: handler = DailyRotatingFileHandler(LOG_DIR, encoding="utf-8") handler.setFormatter(logging.Formatter("%(asctime)s - %(levelname)s - %(message)s", "%Y-%m-%d %H:%M:%S")) logger.addHandler(handler) 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()