data_query_tools.py 41 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000
  1. import hashlib
  2. from zoneinfo import ZoneInfo
  3. from odps import ODPS
  4. from odps.errors import ODPSError
  5. from datetime import date, datetime, timedelta
  6. import json
  7. from pathlib import Path
  8. from examples.demand.mysql import mysql_db
  9. def get_odps_data(sql):
  10. # 配置信息
  11. access_id = 'LTAI9EBa0bd5PrDa'
  12. access_key = 'vAalxds7YxhfOA2yVv8GziCg3Y87v5'
  13. project = 'loghubods'
  14. endpoint = 'http://service.odps.aliyun.com/api'
  15. # 1. 初始化 ODPS 入口
  16. o = ODPS(access_id, access_key, project, endpoint=endpoint)
  17. try:
  18. # 2. 执行 SQL 并获取结果
  19. # execute_sql 会等待任务完成,使用 open_reader 读取数据
  20. with o.execute_sql(sql).open_reader() as reader:
  21. # reader 类似于 Java 中的 List<Record>
  22. # 我们可以直接将其转换为 Python 的 list
  23. records = [record for record in reader]
  24. return records
  25. except ODPSError as e:
  26. print(f"ODPS 错误: {e}")
  27. return None
  28. def execute_odps_sql(sql) -> bool:
  29. # 配置信息
  30. access_id = 'LTAI9EBa0bd5PrDa'
  31. access_key = 'vAalxds7YxhfOA2yVv8GziCg3Y87v5'
  32. project = 'loghubods'
  33. endpoint = 'http://service.odps.aliyun.com/api'
  34. o = ODPS(access_id, access_key, project, endpoint=endpoint)
  35. try:
  36. instance = o.execute_sql(sql)
  37. instance.wait_for_success()
  38. return True
  39. except ODPSError as e:
  40. print(f"ODPS 错误: {e}")
  41. return False
  42. _STRATEGY_GAP = "当下供需gap"
  43. _STRATEGY_GAP_FENCI = "当下供需gap-分词"
  44. _STRATEGY_GAP_ZAOAN = "当下供需gap-早安"
  45. _HIVE_TABLE = "loghubods.dwd_multi_demand_pool_di"
  46. _HIVE_DT_FMT = "%Y%m%d" # 分区格式:yyyymmdd,如 20260519
  47. _CHINA_TZ = ZoneInfo("Asia/Shanghai")
  48. def _hive_partition_dt() -> str:
  49. """中国时区(Asia/Shanghai)当天日期,格式 yyyymmdd。"""
  50. return datetime.now(_CHINA_TZ).date().strftime(_HIVE_DT_FMT)
  51. def _hive_yesterday_partition_dt() -> str:
  52. """中国时区(Asia/Shanghai)昨天日期,格式 yyyymmdd。"""
  53. return (datetime.now(_CHINA_TZ).date() - timedelta(days=1)).strftime(_HIVE_DT_FMT)
  54. def _fenci_demand_filter_key(name: str, merge_leve2: str, type_str: str) -> tuple[str, str, str]:
  55. """当下供需gap-分词策略的去重键(不含 dt,用于跨天过滤)。"""
  56. return (name.strip(), merge_leve2.strip(), type_str.strip())
  57. def _get_demand_filter_words(bizdate: str | None = None) -> list[str]:
  58. """
  59. 从 ODPS demand_filter_word 读取当天过滤词。
  60. dt 为中国时区当天(yyyymmdd);查询失败时返回空列表(不过滤)。
  61. """
  62. dt = bizdate or _hive_partition_dt()
  63. sql = f"""
  64. SELECT filter_word
  65. FROM demand_filter_word
  66. WHERE dt = '{dt}'
  67. """
  68. records = get_odps_data(sql)
  69. if not records:
  70. return []
  71. words: list[str] = []
  72. seen: set[str] = set()
  73. for record in records:
  74. word = str(record[0] or "").strip()
  75. if not word or word in seen:
  76. continue
  77. seen.add(word)
  78. words.append(word)
  79. return words
  80. def _text_contains_filter_word(text: str, filter_words: list[str]) -> bool:
  81. """若 text 包含任一过滤词(等价 LIKE '%word%')返回 True。"""
  82. if not text or not filter_words:
  83. return False
  84. return any(word in text for word in filter_words if word)
  85. def _get_yesterday_fenci_demand_keys() -> set[tuple[str, str, str]]:
  86. """
  87. 从 MySQL demand_content 读取昨天产生的全部需求(含未写入 ODPS 的),
  88. 用于「当下供需gap-分词」策略的跨天去重。
  89. """
  90. yesterday = _hive_yesterday_partition_dt()
  91. rows = mysql_db.select(
  92. "demand_content",
  93. columns="merge_leve2, name, ext_data",
  94. where="dt = %s",
  95. where_params=(yesterday,),
  96. )
  97. keys: set[tuple[str, str, str]] = set()
  98. for row in rows:
  99. name = str(row.get("name") or "").strip()
  100. merge_leve2 = str(row.get("merge_leve2") or "").strip()
  101. if not name or not merge_leve2:
  102. continue
  103. ext_data = _parse_ext_data(row.get("ext_data"))
  104. type_str = str(ext_data.get("type") or "").strip()
  105. keys.add(_fenci_demand_filter_key(name, merge_leve2, type_str))
  106. return keys
  107. def _escape_odps_string(value: object) -> str:
  108. return str(value).replace("'", "''")
  109. def _format_odps_string_array(values: list) -> str:
  110. if not values:
  111. return "ARRAY()"
  112. parts = [f"'{_escape_odps_string(v)}'" for v in values]
  113. return f"ARRAY({','.join(parts)})"
  114. def _parse_ext_data(ext_data_raw: object) -> dict:
  115. if isinstance(ext_data_raw, dict):
  116. return ext_data_raw
  117. if isinstance(ext_data_raw, str) and ext_data_raw.strip():
  118. try:
  119. return json.loads(ext_data_raw)
  120. except json.JSONDecodeError:
  121. return {}
  122. return {}
  123. def _build_hive_select_part(
  124. strategy: str,
  125. demand_id: str,
  126. demand_name: str,
  127. weight: float,
  128. type_str: str,
  129. video_count: int,
  130. video_ids: list[str],
  131. extend_json: str,
  132. ) -> str:
  133. return (
  134. "SELECT "
  135. f"'{_escape_odps_string(strategy)}' AS strategy, "
  136. f"'{_escape_odps_string(demand_id)}' AS demand_id, "
  137. f"'{_escape_odps_string(demand_name)}' AS demand_name, "
  138. f"{weight} AS weight, "
  139. f"'{_escape_odps_string(type_str)}' AS `type`, "
  140. f"{video_count} AS video_count, "
  141. f"{_format_odps_string_array(video_ids)} AS video_list, "
  142. f"'{_escape_odps_string(extend_json)}' AS extend"
  143. )
  144. def _insert_hive_select_parts(select_parts: list[str], partition_dt: str) -> bool:
  145. if not select_parts:
  146. return True
  147. union_sql = "\nUNION ALL\n".join(select_parts)
  148. insert_sql = f"""
  149. INSERT INTO TABLE {_HIVE_TABLE}
  150. PARTITION (dt='{partition_dt}')
  151. (strategy, demand_id, demand_name, weight, `type`, video_count, video_list, extend)
  152. {union_sql}
  153. """
  154. return execute_odps_sql(insert_sql)
  155. def write_dwd_multi_demand_pool_di_to_hive(rows: list[dict]) -> int:
  156. """
  157. 将行数据映射并写入 loghubods.dwd_multi_demand_pool_di(尽力插入,不校验结果)。
  158. 分区与 demand_id 的日期均为中国时区当天(yyyymmdd),不使用行内 dt 字段。
  159. 写入前从 demand_filter_word(dt=当天)拉取过滤词,demand_name 包含任一过滤词则跳过。
  160. 执行两次 INSERT(同表、同分区),策略不同:
  161. 1) 当下供需gap: demand_name=merge_leve2+' '+name, demand_id=md5(strategy+demand_name+type+dt)
  162. 2) 当下供需gap-分词: demand_name=name, demand_id=md5(strategy+name+品类+type+dt)
  163. - 写入前过滤掉昨天 MySQL demand_content 中已产生的全部需求(按 name+品类+type 去重)
  164. """
  165. if not rows:
  166. return 0
  167. china_today = _hive_partition_dt()
  168. filter_words = _get_demand_filter_words(china_today)
  169. yesterday_keys = _get_yesterday_fenci_demand_keys()
  170. gap_parts: list[str] = []
  171. fenci_parts: list[str] = []
  172. fenci_skipped = 0
  173. filter_word_skipped = 0
  174. print(
  175. f"[hive] 过滤词加载: dt={china_today}, count={len(filter_words)}"
  176. )
  177. for row in rows:
  178. merge_leve2 = str(row.get("merge_leve2") or "").strip()
  179. name = str(row.get("name") or "").strip()
  180. if not merge_leve2 or not name:
  181. continue
  182. weight = round(float(row.get("score") or 0.0), 6)
  183. ext_data = _parse_ext_data(row.get("ext_data"))
  184. type_str = str(ext_data.get("type") or "").strip()
  185. video_ids = ext_data.get("video_ids") or []
  186. if not isinstance(video_ids, list):
  187. video_ids = []
  188. video_ids = [str(v).strip() for v in video_ids if v is not None and str(v).strip()]
  189. video_count = len(video_ids)
  190. extend_json = json.dumps({"品类": merge_leve2}, ensure_ascii=False)
  191. demand_name_gap = f"{merge_leve2} {name}"
  192. # 等价 SQL LIKE '%filter_word%':demand_name 包含任一过滤词则不写入 ODPS
  193. if _text_contains_filter_word(demand_name_gap, filter_words):
  194. filter_word_skipped += 1
  195. continue
  196. demand_id_gap = hashlib.md5(
  197. f"{_STRATEGY_GAP}{demand_name_gap}{type_str}{china_today}".encode("utf-8")
  198. ).hexdigest()
  199. gap_parts.append(
  200. _build_hive_select_part(
  201. _STRATEGY_GAP, demand_id_gap, demand_name_gap,
  202. weight, type_str, video_count, video_ids, extend_json,
  203. )
  204. )
  205. filter_key = _fenci_demand_filter_key(name, merge_leve2, type_str)
  206. if filter_key in yesterday_keys:
  207. fenci_skipped += 1
  208. continue
  209. demand_id_fenci = hashlib.md5(
  210. f"{_STRATEGY_GAP_FENCI}{name}{merge_leve2}{type_str}{china_today}".encode("utf-8")
  211. ).hexdigest()
  212. fenci_parts.append(
  213. _build_hive_select_part(
  214. _STRATEGY_GAP_FENCI, demand_id_fenci, name,
  215. weight, type_str, video_count, video_ids, extend_json,
  216. )
  217. )
  218. print(
  219. f"[hive] 过滤词跳过: skipped={filter_word_skipped}, "
  220. f"gap_to_write={len(gap_parts)}, fenci_to_write={len(fenci_parts)}"
  221. )
  222. if not gap_parts:
  223. return 0
  224. print(
  225. f"[hive] {_STRATEGY_GAP_FENCI} 过滤昨天需求: "
  226. f"yesterday_dt={_hive_yesterday_partition_dt()}, "
  227. f"yesterday_total={len(yesterday_keys)}, skipped={fenci_skipped}, "
  228. f"to_write={len(fenci_parts)}"
  229. )
  230. _insert_hive_select_parts(gap_parts, china_today)
  231. if fenci_parts:
  232. _insert_hive_select_parts(fenci_parts, china_today)
  233. return len(gap_parts) + len(fenci_parts)
  234. def write_dwd_zaoan_demand_pool_to_hive(rows: list[dict]) -> int:
  235. """
  236. 写入策略「当下供需gap-早安」到 loghubods.dwd_multi_demand_pool_di。
  237. demand_name 只用需求词本身,不拼品类。
  238. demand_id = md5(strategy + demand_name + type + dt)
  239. 该策略已由上游筛过,不再套用过滤词 / 昨天分词去重。
  240. """
  241. if not rows:
  242. return 0
  243. china_today = _hive_partition_dt()
  244. select_parts: list[str] = []
  245. for row in rows:
  246. merge_leve2 = str(row.get("merge_leve2") or "").strip()
  247. name = str(row.get("name") or "").strip()
  248. if not merge_leve2 or not name:
  249. continue
  250. weight = round(float(row.get("score") or 0.0), 6)
  251. ext_data = _parse_ext_data(row.get("ext_data"))
  252. type_str = str(ext_data.get("type") or "").strip()
  253. video_ids = ext_data.get("video_ids") or []
  254. if not isinstance(video_ids, list):
  255. video_ids = []
  256. video_ids = [str(v).strip() for v in video_ids if v is not None and str(v).strip()]
  257. video_count = len(video_ids)
  258. extend_json = json.dumps({"品类": merge_leve2}, ensure_ascii=False)
  259. demand_name = name
  260. demand_id = hashlib.md5(
  261. f"{_STRATEGY_GAP_ZAOAN}{demand_name}{type_str}{china_today}".encode("utf-8")
  262. ).hexdigest()
  263. select_parts.append(
  264. _build_hive_select_part(
  265. _STRATEGY_GAP_ZAOAN, demand_id, demand_name,
  266. weight, type_str, video_count, video_ids, extend_json,
  267. )
  268. )
  269. if not select_parts:
  270. print(f"[hive] {_STRATEGY_GAP_ZAOAN} 无有效行,跳过")
  271. return 0
  272. print(f"[hive] {_STRATEGY_GAP_ZAOAN} 写入: dt={china_today}, rows={len(select_parts)}")
  273. ok = _insert_hive_select_parts(select_parts, china_today)
  274. if not ok:
  275. raise RuntimeError(f"写入 Hive 失败: strategy={_STRATEGY_GAP_ZAOAN}, dt={china_today}")
  276. return len(select_parts)
  277. def write_feature_point_data_to_hive(names: list[str]) -> int:
  278. """
  279. 将需求名称写入 Hive 表 feature_point_data(按北京时间当天分区)。
  280. 仅写入以下字段:
  281. - 特征点
  282. - 总分发曝光pv(固定 5000)
  283. - 质bn_rovn(固定 0.1)
  284. """
  285. normalized_names = [str(name).strip() for name in names if name is not None and str(name).strip()]
  286. if not normalized_names:
  287. return 0
  288. dt = datetime.now(ZoneInfo("Asia/Shanghai")).strftime("%Y%m%d")
  289. select_parts = []
  290. for name in normalized_names:
  291. safe_name = name.replace("'", "''")
  292. select_parts.append(
  293. "SELECT "
  294. f"'{safe_name}' AS `特征点`, "
  295. "5000 AS `总分发曝光pv`, "
  296. "0.1 AS `质bn_rovn`"
  297. )
  298. union_sql = "\nUNION ALL\n".join(select_parts)
  299. insert_sql = f"""
  300. INSERT INTO TABLE feature_point_data
  301. PARTITION (dt='{dt}')
  302. (`特征点`, `总分发曝光pv`, `质bn_rovn`)
  303. {union_sql}
  304. """
  305. ok = execute_odps_sql(insert_sql)
  306. if not ok:
  307. return 0
  308. return len(normalized_names)
  309. def get_demand_merge_level2_names():
  310. date_time = datetime.now(ZoneInfo("Asia/Shanghai")).date() - timedelta(days=1)
  311. day = date_time.strftime("%Y%m%d")
  312. count = 50
  313. sql_query = f'''
  314. select *
  315. from (
  316. select
  317. dt,
  318. merge二级品类,
  319. sum(当日分发曝光pv) as 分发曝光pv,
  320. sum(累计分享回流uv) AS bn_总回流,
  321. sum(当日分发回流uv)/(sum(当日分发曝光pv)+100) as 质bn_rovn,
  322. case when sum(当日分发曝光pv)>=10000 then
  323. case when sum(当日分发回流uv)/(sum(当日分发曝光pv)+100)<0.035
  324. then -1*(count(DISTINCT 视频id)/avg(总日分发视频数))/((sum(累计分享回流uv)/avg(总日回流uv)))
  325. else 10*(sum(累计分享回流uv)/avg(总日回流uv)*sum(当日分发回流uv)/(sum(当日分发曝光pv)+100))/(count(DISTINCT 视频id)/avg(总日分发视频数))
  326. end
  327. else 0 end AS 总供需分,
  328. case when sum(当日分发曝光pv)>=10000 then
  329. case when sum(当日分发回流uv)/(sum(当日分发曝光pv)+100)<0.035
  330. then -1*(COUNT(DISTINCT CASE WHEN 推荐天数间隔<3 THEN 视频id END ) /avg(总日分发视频数))/(sum(累计分享回流uv)/avg(总日回流uv))
  331. else 10*(sum(累计分享回流uv)/avg(总日回流uv)*sum(当日分发回流uv)/(sum(当日分发曝光pv)+1000))/(COUNT(DISTINCT CASE WHEN 推荐天数间隔<3 THEN 视频id END ) /avg(总日分发视频数))
  332. end
  333. else 0 end AS 新供需分,
  334. count(DISTINCT 视频id) as 分发视频量,
  335. count(DISTINCT if(推荐天数间隔<3,视频id,null)) as 3日新推荐视频量,
  336. case when sum(当日分发曝光pv)>=10000 and sum(当日分发回流uv)/(sum(当日分发曝光pv)+100)>0.035
  337. then (avg(总日分发视频数)*(10*(sum(当日分发回流uv)/(sum(当日分发曝光pv)+100))*(sum(累计分享回流uv)/avg(总日回流uv) ))/0.5-count(DISTINCT 视频id))/3
  338. end as 缺量,
  339. case when sum(当日分发曝光pv)>=10000 and sum(当日分发回流uv)/(sum(当日分发曝光pv)+100)<=0.035
  340. then (avg(总日分发视频数)*(10*(sum(当日分发回流uv)/(sum(当日分发曝光pv)+100))*(sum(累计分享回流uv)/avg(总日回流uv) ))/(2)-count(DISTINCT 视频id))/3
  341. end as 控量,
  342. avg(总日回流uv) AS 总日回流uv,
  343. avg(总日分发视频数) AS 总日分发视频数,
  344. avg(总日推荐视频数) AS 总日推荐视频数,
  345. COUNT(DISTINCT CASE WHEN 总回流uv>0 THEN 视频id END )/avg(总日分发视频数) AS 回流视频个数占比,
  346. sum(当日分发回流uv) AS bn_当日分发回流,
  347. sum(当日分发回流uv)/avg(总日回流uv) AS 分发拉回回流uv占比,
  348. sum(累计分享回流uv)/avg(总日回流uv) AS 回流uv占比,
  349. count(DISTINCT 视频id)/avg(总日分发视频数) AS 分发视频量占比,
  350. COUNT(DISTINCT CASE WHEN 是否当日新推荐=1 THEN 视频id END ) /avg(总日分发视频数) AS 新推荐视频量占比
  351. from loghubods.video_dimension_detail_add_column
  352. where dt = '{day}'
  353. group by dt, merge二级品类
  354. ) t1
  355. where t1.缺量>= {count}
  356. '''
  357. data = get_odps_data(sql_query)
  358. result_list = []
  359. if data:
  360. for r in data:
  361. lack_count = r[9]
  362. if lack_count > 1000:
  363. count = 70
  364. elif 500 < lack_count <= 1000:
  365. count = 60
  366. elif 100 < lack_count <= 500:
  367. count = 40
  368. elif 50 < lack_count <= 100:
  369. count = 20
  370. else:
  371. count = 10
  372. if count == 0:
  373. continue
  374. result_list.append({
  375. "cluster_name": r[1],
  376. "platform_type": "piaoquan",
  377. "count": count,
  378. })
  379. return result_list
  380. def get_rov_by_merge_leve2_and_video_ids(merge_leve2, video_ids):
  381. merge_level_in_clause = f"'{merge_leve2}'"
  382. video_ids_in_clause = ", ".join([f"'{video_id}'" for video_id in video_ids])
  383. end_date = (date.today() - timedelta(days=1)).strftime("%Y%m%d")
  384. start_date = (date.today() - timedelta(days=14)).strftime("%Y%m%d")
  385. # sql_query = f'''
  386. # SELECT
  387. # v.videoid,
  388. # CASE
  389. # WHEN COALESCE(SUM(COALESCE(t3.`当日分发曝光pv`, 0)), 0) < 1000 THEN 0
  390. # ELSE COALESCE(AVG(NULLIF(t3.rov_t0, 0)), 0)
  391. # END AS avg_rov_t0
  392. # FROM
  393. # (
  394. # SELECT
  395. # t2.videoid,
  396. # t2.merge_leve2
  397. # FROM videoods.content_profile t1
  398. # JOIN loghubods.video_merge_tag t2
  399. # ON t1.content_id = t2.videoid
  400. # WHERE
  401. # t1.status = 3
  402. # AND t1.is_deleted = 0
  403. # AND t2.merge_leve2 IN ({merge_level_in_clause})
  404. # ) v
  405. # LEFT JOIN loghubods.video_dimension_detail_add_column t3
  406. # ON v.videoid = t3.视频id
  407. # AND t3.dt >= '{start_date}'
  408. # AND t3.dt <= '{end_date}'
  409. # WHERE v.videoid in ({video_ids_in_clause})
  410. # GROUP BY
  411. # v.videoid
  412. # ;
  413. # '''
  414. sql_query = f'''
  415. SELECT
  416. CAST(t3.视频id AS STRING) AS 视频id_str,
  417. CASE
  418. WHEN COALESCE(SUM(COALESCE(t3.`当日分发曝光pv`, 0)), 0) < 1000 THEN 0
  419. ELSE COALESCE(AVG(NULLIF(t3.rov_t0, 0)), 0)
  420. END AS avg_rov_t0
  421. FROM
  422. loghubods.video_dimension_detail_add_column t3
  423. WHERE t3.视频id in ({video_ids_in_clause})
  424. AND t3.dt >= '{start_date}'
  425. AND t3.dt <= '{end_date}'
  426. GROUP BY
  427. t3.视频id
  428. ;
  429. '''
  430. data = get_odps_data(sql_query)
  431. result_dict = {}
  432. if data:
  433. result_dict = {r[0]: r[1] for r in data}
  434. return result_dict
  435. def get_rov_by_tree_and_video_ids(video_ids):
  436. video_ids_in_clause = ", ".join([f"'{video_id}'" for video_id in video_ids])
  437. last_year_today = date.today() - timedelta(days=365)
  438. start_date = last_year_today.strftime("%Y%m%d")
  439. end_date = (last_year_today + timedelta(days=7)).strftime("%Y%m%d")
  440. sql_query = f'''
  441. SELECT
  442. CAST(t3.视频id AS STRING) AS 视频id_str,
  443. CASE
  444. WHEN COALESCE(SUM(COALESCE(t3.`当日分发曝光pv`, 0)), 0) < 1000 THEN 0
  445. ELSE COALESCE(AVG(NULLIF(t3.rov_t0, 0)), 0)
  446. END AS avg_rov_t0
  447. FROM
  448. loghubods.video_dimension_detail_add_column t3
  449. WHERE t3.视频id in ({video_ids_in_clause})
  450. AND t3.dt >= '{start_date}'
  451. AND t3.dt <= '{end_date}'
  452. GROUP BY
  453. t3.视频id
  454. ;
  455. '''
  456. data = get_odps_data(sql_query)
  457. result_dict = {}
  458. if data:
  459. result_dict = {r[0]: r[1] for r in data}
  460. return result_dict
  461. def get_changwen_weight(account_name):
  462. bizdatemax_date = date.today() - timedelta(days=1)
  463. bizdatemin_date = bizdatemax_date - timedelta(days=30)
  464. bizdatemax = bizdatemax_date.strftime("%Y%m%d")
  465. bizdatemin = bizdatemin_date.strftime("%Y%m%d")
  466. sql_query = f'''
  467. SELECT
  468. 公众号名
  469. ,videoid
  470. ,一级品类
  471. ,二级品类
  472. ,头部曝光
  473. ,头部曝光uv
  474. ,头部realplay
  475. ,头部realplay_uv
  476. ,头部分享
  477. ,头部分享uv
  478. ,头部回流人数 AS 头部回流数
  479. ,推荐曝光数
  480. ,当日分发曝光uv
  481. ,推荐realplay
  482. ,分发realplay_uv
  483. ,推荐分享数
  484. ,当日分发分享uv
  485. ,推荐回流数
  486. ,当日回流进入分发曝光次数 AS vov分子
  487. FROM (
  488. SELECT DISTINCT a.公众号名
  489. ,a.videoid
  490. ,e.merge_leve1 AS 一级品类
  491. ,e.merge_leve2 AS 二级品类
  492. ,a.title
  493. ,a.进入分发人数
  494. ,头部曝光pv AS 头部曝光
  495. ,头部realplay_pv AS 头部realplay
  496. ,头部分享pv AS 头部分享
  497. ,a.当日分发曝光pv AS 推荐曝光数
  498. ,a.当日分发播放pv
  499. ,分发realplay_pv AS 推荐realplay
  500. ,分发realplay_pv / a.当日分发播放pv AS 真实播放率pv
  501. ,当日分发播放uv
  502. ,c.realplay_uv AS 分发真实播uv
  503. ,c.realplay_uv / a.当日分发播放uv AS 真实播放率uv
  504. ,a.当日分发分享pv AS 推荐分享数
  505. ,a.当日分发分享pv / a.当日分发曝光pv AS str
  506. ,NVL(b.当日分发回流人数,0) AS 推荐回流数
  507. ,NVL(b.当日回流进入分发人数,0) AS 当日回流进入分发人数
  508. ,NVL(b.当日回流进入分发曝光次数,0) AS 当日回流进入分发曝光次数
  509. ,NVL(b.当日回流进入分发曝光次数,0) / a.当日分发曝光pv AS vov分子
  510. ,d.头部回流人数
  511. ,当日分发曝光uv
  512. ,头部曝光uv
  513. ,当日分发分享uv
  514. ,头部分享uv
  515. ,分发realplay_uv
  516. ,头部realplay_uv
  517. FROM (
  518. SELECT account_name AS 公众号名
  519. ,videoid
  520. ,title
  521. ,COUNT(DISTINCT mid) AS 进入分发人数
  522. ,COUNT(
  523. CASE WHEN pagesource REGEXP 'category$|recommend$|-pages/user-videos-detail$' AND businesstype = 'videoView' THEN mid END
  524. ) AS 当日分发曝光pv
  525. ,COUNT(DISTINCT
  526. CASE WHEN pagesource REGEXP 'category$|recommend$|-pages/user-videos-detail$' AND businesstype = 'videoView' THEN mid END
  527. ) AS 当日分发曝光uv
  528. ,COUNT(
  529. CASE WHEN pagesource REGEXP 'pages/user-videos-share$' AND businesstype = 'videoView' THEN mid END
  530. ) AS 头部曝光pv
  531. ,COUNT(DISTINCT
  532. CASE WHEN pagesource REGEXP 'pages/user-videos-share$' AND businesstype = 'videoView' THEN mid END
  533. ) AS 头部曝光uv
  534. ,COUNT(
  535. CASE WHEN pagesource REGEXP 'category$|recommend$|-pages/user-videos-detail$' AND businesstype = 'videoPlay' THEN mid END
  536. ) AS 当日分发播放pv
  537. ,COUNT(DISTINCT
  538. CASE WHEN pagesource REGEXP 'category$|recommend$|-pages/user-videos-detail$' AND businesstype = 'videoPlay' THEN mid END
  539. ) AS 当日分发播放uv
  540. ,COUNT(
  541. CASE WHEN pagesource REGEXP 'category$|recommend$|-pages/user-videos-detail$' AND businesstype = 'videoShareFriend' THEN mid END
  542. ) AS 当日分发分享pv
  543. ,COUNT(DISTINCT
  544. CASE WHEN pagesource REGEXP 'category$|recommend$|-pages/user-videos-detail$' AND businesstype = 'videoShareFriend' THEN mid END
  545. ) AS 当日分发分享uv
  546. ,COUNT(
  547. CASE WHEN pagesource REGEXP 'pages/user-videos-share$' AND businesstype = 'videoShareFriend' THEN mid END
  548. ) AS 头部分享pv
  549. ,COUNT(DISTINCT
  550. CASE WHEN pagesource REGEXP 'pages/user-videos-share$' AND businesstype = 'videoShareFriend' THEN mid END
  551. ) AS 头部分享uv
  552. FROM (
  553. SELECT DISTINCT a.mid
  554. ,a.videoid
  555. ,a.businesstype
  556. ,a.pagesource
  557. ,a.subsessionid
  558. ,account_name
  559. ,e.title
  560. FROM loghubods.video_action_log_rp a
  561. LEFT JOIN loghubods.user_wechat_identity_info_ha b
  562. ON a.mid = CONCAT('weixin_openid_',b.open_id)
  563. AND b.dt = MAX_PT("loghubods.user_wechat_identity_info_ha")
  564. LEFT JOIN loghubods.gzh_fans_info d
  565. ON b.union_id = d.union_id
  566. AND d.dt = MAX_PT("loghubods.gzh_fans_info")
  567. LEFT JOIN videoods.wx_video e
  568. ON a.videoid = e.id
  569. WHERE a.dt >= '{bizdatemin}'
  570. AND a.dt <= '{bizdatemax}'
  571. AND businesstype IN ('videoView','videoPlay','videoShareFriend')
  572. AND d.user_create_time IS NOT NULL
  573. AND account_name = '{account_name}'
  574. AND a.videoid IN (
  575. SELECT
  576. DISTINCT content_id AS videoid
  577. FROM
  578. videoods.content_profile
  579. WHERE status=3
  580. AND is_deleted = 0
  581. )
  582. ) t
  583. GROUP BY 公众号名
  584. ,videoid
  585. ,title
  586. ) a
  587. LEFT JOIN (
  588. SELECT t.account_name AS 公众号名
  589. ,t.videoid
  590. ,COUNT(DISTINCT s.machinecode) AS 当日分发回流人数
  591. ,COUNT(DISTINCT v.mid) AS 当日回流进入分发人数
  592. ,COUNT(v.mid) AS 当日回流进入分发曝光次数
  593. FROM (
  594. SELECT DISTINCT a.subsessionid
  595. ,a.videoid
  596. ,a.mid
  597. ,d.account_name
  598. ,GET_JSON_OBJECT(extparams,'$.recomTraceId') AS recomtraceid
  599. FROM loghubods.video_action_log_rp a
  600. LEFT JOIN loghubods.user_wechat_identity_info_ha b
  601. ON a.mid = CONCAT('weixin_openid_',b.open_id)
  602. AND b.dt = MAX_PT("loghubods.user_wechat_identity_info_ha")
  603. LEFT JOIN loghubods.gzh_fans_info d
  604. ON b.union_id = d.union_id
  605. AND d.dt = MAX_PT("loghubods.gzh_fans_info")
  606. WHERE a.dt >= '{bizdatemin}'
  607. AND a.dt <= '{bizdatemax}'
  608. AND a.businesstype = 'videoShareFriend'
  609. AND a.pagesource REGEXP 'category$|recommend$|-pages/user-videos-detail$'
  610. AND d.user_create_time IS NOT NULL
  611. AND d.account_name = '{account_name}'
  612. ) t
  613. LEFT JOIN (
  614. SELECT DISTINCT subsessionid
  615. ,machinecode
  616. ,recomtraceid
  617. ,clickobjectid
  618. FROM loghubods.user_share_log
  619. WHERE dt >= '{bizdatemin}'
  620. AND dt <= '{bizdatemax}'
  621. AND topic = 'click'
  622. ) s
  623. ON t.recomtraceid = s.recomtraceid
  624. AND t.videoid = s.clickobjectid
  625. LEFT JOIN (
  626. SELECT subsessionid
  627. ,mid
  628. ,videoid
  629. FROM loghubods.video_action_log_rp
  630. WHERE dt >= '{bizdatemin}'
  631. AND dt <= '{bizdatemax}'
  632. AND pagesource REGEXP 'category$|recommend$|-pages/user-videos-detail$'
  633. AND businesstype = 'videoView'
  634. ) v
  635. ON s.subsessionid = v.subsessionid
  636. AND s.machinecode = v.mid
  637. GROUP BY account_name
  638. ,t.videoid
  639. ) b
  640. ON a.公众号名 = b.公众号名
  641. AND a.videoid = b.videoid
  642. LEFT JOIN (
  643. SELECT d.account_name AS 公众号名
  644. ,a.videoid
  645. ,COUNT(DISTINCT a.mid) AS realplay_uv
  646. ,COUNT(
  647. CASE WHEN a.pagesource REGEXP 'category$|recommend$|-pages/user-videos-detail$' THEN a.mid END
  648. ) AS 分发realplay_pv
  649. ,COUNT(CASE WHEN a.pagesource REGEXP 'pages/user-videos-share$' THEN a.mid END) AS 头部realplay_pv
  650. ,COUNT(DISTINCT
  651. CASE WHEN a.pagesource REGEXP 'category$|recommend$|-pages/user-videos-detail$' THEN a.mid END
  652. ) AS 分发realplay_uv
  653. ,COUNT(DISTINCT CASE WHEN a.pagesource REGEXP 'pages/user-videos-share$' THEN a.mid END) AS 头部realplay_uv
  654. FROM loghubods.ods_video_play_log_day a
  655. LEFT JOIN (
  656. SELECT DISTINCT open_id
  657. ,union_id
  658. FROM loghubods.user_wechat_identity_info_ha
  659. WHERE dt = MAX_PT("loghubods.user_wechat_identity_info_ha")
  660. ) b
  661. ON a.mid = CONCAT('weixin_openid_',b.open_id)
  662. LEFT JOIN loghubods.gzh_fans_info d
  663. ON b.union_id = d.union_id
  664. AND d.dt = MAX_PT("loghubods.gzh_fans_info")
  665. WHERE a.dt >= '{bizdatemin}'
  666. AND a.dt <= '{bizdatemax}'
  667. AND a.businesstype = 'videoRealPlay'
  668. AND d.user_create_time IS NOT NULL
  669. AND d.account_name = '{account_name}'
  670. GROUP BY d.account_name
  671. ,a.videoid
  672. ORDER BY 分发realplay_pv DESC
  673. ) c
  674. ON a.公众号名 = c.公众号名
  675. AND a.videoid = c.videoid
  676. LEFT JOIN (
  677. SELECT t.account_name AS 公众号名
  678. ,t.videoid
  679. ,COUNT(DISTINCT s.machinecode) AS 头部回流人数
  680. FROM (
  681. SELECT DISTINCT a.shareobjectid AS videoid
  682. ,a.shareid
  683. ,a.machinecode
  684. ,d.account_name
  685. FROM loghubods.user_share_log a
  686. LEFT JOIN loghubods.user_wechat_identity_info_ha b
  687. ON a.machinecode = CONCAT('weixin_openid_',b.open_id)
  688. AND b.dt = MAX_PT("loghubods.user_wechat_identity_info_ha")
  689. LEFT JOIN loghubods.gzh_fans_info d
  690. ON b.union_id = d.union_id
  691. AND d.dt = MAX_PT("loghubods.gzh_fans_info")
  692. WHERE a.dt >= '{bizdatemin}'
  693. AND a.dt <= '{bizdatemax}'
  694. AND a.topic = 'share'
  695. AND a.pagesource REGEXP 'pages/user-videos-share$'
  696. AND d.user_create_time IS NOT NULL
  697. AND d.account_name = '{account_name}'
  698. ) t
  699. LEFT JOIN (
  700. SELECT DISTINCT shareid
  701. ,machinecode
  702. ,clickobjectid
  703. FROM loghubods.user_share_log
  704. WHERE dt >= '{bizdatemin}'
  705. AND dt <= '{bizdatemax}'
  706. AND topic = 'click'
  707. ) s
  708. ON t.shareid = s.shareid
  709. GROUP BY account_name
  710. ,t.videoid
  711. ) d
  712. ON a.公众号名 = d.公众号名
  713. AND a.videoid = d.videoid
  714. LEFT JOIN loghubods.video_merge_tag e
  715. ON a.videoid = e.videoid
  716. )
  717. ORDER BY 推荐曝光数 DESC
  718. '''
  719. result_list = []
  720. data = get_odps_data(sql_query)
  721. if data:
  722. for r in data:
  723. result_list.append(
  724. {
  725. "account_name": r[0],
  726. "videoid": r[1],
  727. "一级品类": r[2],
  728. "二级品类": r[3],
  729. "ext_data": {
  730. "头部曝光": r[4],
  731. "头部曝光uv": r[5],
  732. "头部realplay": r[6],
  733. "头部realplay_uv": r[7],
  734. "头部分享": r[8],
  735. "头部分享uv": r[9],
  736. "头部回流数": r[10],
  737. "推荐曝光数": r[11],
  738. "当日分发曝光uv": r[12],
  739. "推荐realplay": r[13],
  740. "分发realplay_uv": r[14],
  741. "推荐分享数": r[15],
  742. "当日分发分享uv": r[16],
  743. "推荐回流数": r[17],
  744. "vov分子": r[18],
  745. },
  746. }
  747. )
  748. # 输出到 examples/demand/data/changwen_data/
  749. output_dir = Path(__file__).parent / "data" / "changwen_data"
  750. output_dir.mkdir(parents=True, exist_ok=True)
  751. output_file = output_dir / f"{account_name}.json"
  752. with output_file.open("w", encoding="utf-8") as f:
  753. json.dump(result_list, f, ensure_ascii=False, indent=2)
  754. return result_list
  755. def get_zengzhang_weight(account_name):
  756. bizdatemax_date = date.today() - timedelta(days=1)
  757. bizdatemin_date = bizdatemax_date - timedelta(days=30)
  758. bizdatemax = bizdatemax_date.strftime("%Y%m%d")
  759. bizdatemin = bizdatemin_date.strftime("%Y%m%d")
  760. sql_query = f'''
  761. SELECT 合作方名
  762. ,合作方简称
  763. ,videoid
  764. ,一级品类
  765. ,二级品类
  766. ,SUM(头部曝光) as 头部曝光
  767. ,SUM(头部曝光uv) as 头部曝光uv
  768. ,SUM(头部realplay) as 头部realplay
  769. ,SUM(头部realplay_uv) as 头部realplay_uv
  770. ,SUM(头部分享) as 头部分享
  771. ,SUM(头部分享uv) as 头部分享uv
  772. ,SUM(头部回流数) as 头部回流数
  773. ,SUM(推荐曝光数) as 推荐曝光数
  774. ,SUM(当日分发曝光uv) as 当日分发曝光uv
  775. ,SUM(推荐realplay) as 推荐realplay
  776. ,SUM(分发realplay_uv) as 分发realplay_uv
  777. ,SUM(推荐分享数) as 推荐分享数
  778. ,SUM(当日分发分享uv) as 当日分发分享uv
  779. ,SUM(推荐回流数) as 推荐回流数
  780. ,SUM(vov分子) as vov分子
  781. FROM loghubods.dws_growth_partner_vid_data
  782. WHERE dt BETWEEN '{bizdatemin}' AND '{bizdatemax}'
  783. AND 合作方名 = '{account_name}'
  784. GROUP BY 合作方名
  785. ,合作方简称
  786. ,videoid
  787. ,一级品类
  788. ,二级品类
  789. ORDER BY SUM(推荐曝光数)
  790. ;
  791. '''
  792. result_list = []
  793. data = get_odps_data(sql_query)
  794. if data:
  795. for r in data:
  796. result_list.append(
  797. {
  798. "account_name": r[0],
  799. "合作方简称": r[1],
  800. "videoid": r[2],
  801. "一级品类": r[3],
  802. "二级品类": r[4],
  803. "ext_data": {
  804. "头部曝光": r[5],
  805. "头部曝光uv": r[6],
  806. "头部realplay": r[7],
  807. "头部realplay_uv": r[8],
  808. "头部分享": r[9],
  809. "头部分享uv": r[10],
  810. "头部回流数": r[11],
  811. "推荐曝光数": r[12],
  812. "当日分发曝光uv": r[13],
  813. "推荐realplay": r[14],
  814. "分发realplay_uv": r[15],
  815. "推荐分享数": r[16],
  816. "当日分发分享uv": r[17],
  817. "推荐回流数": r[18],
  818. "vov分子": r[19],
  819. },
  820. }
  821. )
  822. # 输出到 examples/demand/data/zengzhang_data/
  823. output_dir = Path(__file__).parent / "data" / "zengzhang_data"
  824. output_dir.mkdir(parents=True, exist_ok=True)
  825. output_file = output_dir / f"{account_name}.json"
  826. with output_file.open("w", encoding="utf-8") as f:
  827. json.dump(result_list, f, ensure_ascii=False, indent=2)
  828. return result_list
  829. def get_merge_leve2_by_video_ids(video_ids, batch_size=2000):
  830. result = {}
  831. if not video_ids:
  832. return result
  833. normalized_ids = [str(video_id) for video_id in video_ids if video_id is not None]
  834. for i in range(0, len(normalized_ids), batch_size):
  835. batch_ids = normalized_ids[i:i + batch_size]
  836. escaped_ids = [video_id.replace("'", "''") for video_id in batch_ids]
  837. video_ids_in_clause = ", ".join([f"'{video_id}'" for video_id in escaped_ids])
  838. sql_query = f'''
  839. SELECT videoid, merge_leve2
  840. FROM loghubods.video_merge_tag
  841. WHERE videoid IN ({video_ids_in_clause})
  842. '''
  843. data = get_odps_data(sql_query)
  844. if not data:
  845. continue
  846. for row in data:
  847. result[str(row[0])] = row[1]
  848. return result
  849. def get_all_decode_task_result_rows():
  850. return mysql_db.select(
  851. "workflow_decode_task_result",
  852. columns="id, channel_content_id, merge_leve2",
  853. )
  854. def update_decode_task_result_merge_leve2(channel_content_id, merge_leve2):
  855. return mysql_db.update(
  856. "workflow_decode_task_result",
  857. {"merge_leve2": str(merge_leve2)},
  858. "channel_content_id = %s",
  859. (str(channel_content_id),),
  860. )
  861. def backfill_merge_leve2_for_decode_task_result():
  862. rows = get_all_decode_task_result_rows()
  863. updated_count = 0
  864. skipped_count = 0
  865. valid_content_ids = []
  866. for row in rows:
  867. channel_content_id = row.get("channel_content_id")
  868. if channel_content_id is None:
  869. skipped_count += 1
  870. continue
  871. channel_content_id = str(channel_content_id)
  872. if len(channel_content_id) > 8:
  873. skipped_count += 1
  874. continue
  875. valid_content_ids.append(channel_content_id)
  876. merge_leve2_map = get_merge_leve2_by_video_ids(valid_content_ids, batch_size=2000)
  877. for channel_content_id in valid_content_ids:
  878. merge_leve2 = merge_leve2_map.get(channel_content_id)
  879. if not merge_leve2:
  880. continue
  881. affected = update_decode_task_result_merge_leve2(channel_content_id, merge_leve2)
  882. if affected > 0:
  883. updated_count += affected
  884. return {
  885. "total": len(rows),
  886. "updated": updated_count,
  887. "skipped": skipped_count,
  888. }
  889. #
  890. # if __name__ == '__main__':
  891. # backfill_merge_leve2_for_decode_task_result()