data_query_tools.py 41 KB

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