data_query_tools.py 37 KB

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