data_query_tools.py 38 KB

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