yumeng 2 روز پیش
والد
کامیت
74dd29b144
1فایلهای تغییر یافته به همراه32 افزوده شده و 3 حذف شده
  1. 32 3
      scripts/honor_backfill_oaid.py

+ 32 - 3
scripts/honor_backfill_oaid.py

@@ -48,6 +48,12 @@ def parse_args() -> argparse.Namespace:
     parser.add_argument("--redis-db", type=int, default=5)
     parser.add_argument("--redis-db", type=int, default=5)
     parser.add_argument("--redis-prefix", default="adx:honor:")
     parser.add_argument("--redis-prefix", default="adx:honor:")
     parser.add_argument("--limit", type=int, default=0, help="Only process first N extracted mappings")
     parser.add_argument("--limit", type=int, default=0, help="Only process first N extracted mappings")
+    parser.add_argument(
+        "--progress-every",
+        type=int,
+        default=100000,
+        help="Print progress every N scanned log lines / processed mappings",
+    )
     parser.add_argument("--dry-run", action="store_true")
     parser.add_argument("--dry-run", action="store_true")
     return parser.parse_args()
     return parser.parse_args()
 
 
@@ -72,11 +78,19 @@ def expand_logs(patterns: Iterable[str]) -> List[Path]:
     return deduped
     return deduped
 
 
 
 
-def scan_logs(paths: Iterable[Path]) -> Dict[str, LogEntry]:
+def scan_logs(paths: Iterable[Path], progress_every: int) -> Dict[str, LogEntry]:
     result: Dict[str, LogEntry] = {}
     result: Dict[str, LogEntry] = {}
     for path in paths:
     for path in paths:
+        print(f"[scan] start file={path}")
+        line_count = 0
+        match_count = 0
         with path.open("r", encoding="utf-8", errors="ignore") as fh:
         with path.open("r", encoding="utf-8", errors="ignore") as fh:
             for line in fh:
             for line in fh:
+                line_count += 1
+                if progress_every > 0 and line_count % progress_every == 0:
+                    print(
+                        f"[scan] file={path} lines={line_count} mappings={len(result)} fileMatches={match_count}"
+                    )
                 if "[荣耀][曝光][入参]" not in line and "[荣耀][点击][入参]" not in line:
                 if "[荣耀][曝光][入参]" not in line and "[荣耀][点击][入参]" not in line:
                     continue
                     continue
                 trace_match = TRACE_ID_RE.search(line)
                 trace_match = TRACE_ID_RE.search(line)
@@ -92,6 +106,8 @@ def scan_logs(paths: Iterable[Path]) -> Dict[str, LogEntry]:
                 if not oaid:
                 if not oaid:
                     continue
                     continue
                 result[trace_id] = LogEntry(trace_id=trace_id, oaid=oaid)
                 result[trace_id] = LogEntry(trace_id=trace_id, oaid=oaid)
+                match_count += 1
+        print(f"[scan] done file={path} lines={line_count} mappings={len(result)} fileMatches={match_count}")
     return result
     return result
 
 
 
 
@@ -220,7 +236,13 @@ def main() -> int:
         print("no log files matched")
         print("no log files matched")
         return 1
         return 1
 
 
-    mappings = scan_logs(log_files)
+    print(
+        f"[start] logs={len(log_files)} dry_run={args.dry_run} limit={args.limit or 'all'} progress_every={args.progress_every}"
+    )
+    for path in log_files:
+        print(f"[start] log={path}")
+
+    mappings = scan_logs(log_files, args.progress_every)
     if not mappings:
     if not mappings:
         print("no honor traceId/oaid mappings found")
         print("no honor traceId/oaid mappings found")
         return 1
         return 1
@@ -228,6 +250,7 @@ def main() -> int:
     if args.limit > 0:
     if args.limit > 0:
         items = items[:args.limit]
         items = items[:args.limit]
     trace_ids = [trace_id for trace_id, _ in items]
     trace_ids = [trace_id for trace_id, _ in items]
+    print(f"[scan] extracted_mappings={len(items)}")
 
 
     conn = pymysql.connect(
     conn = pymysql.connect(
         host=args.mysql_host,
         host=args.mysql_host,
@@ -259,7 +282,12 @@ def main() -> int:
             "redis_missing": 0,
             "redis_missing": 0,
         }
         }
 
 
-        for trace_id, entry in items:
+        for index, (trace_id, entry) in enumerate(items, start=1):
+            if args.progress_every > 0 and index % args.progress_every == 0:
+                print(
+                    f"[apply] processed={index}/{len(items)} mysql_updated={stats['mysql_updated']} "
+                    f"redis_updated={stats['redis_keys_updated']} db_missing={stats['db_missing']}"
+                )
             row = db_rows.get(trace_id)
             row = db_rows.get(trace_id)
             if not row:
             if not row:
                 stats["db_missing"] += 1
                 stats["db_missing"] += 1
@@ -305,6 +333,7 @@ def main() -> int:
             conn.rollback()
             conn.rollback()
         else:
         else:
             conn.commit()
             conn.commit()
+        print(f"[done] processed={len(items)}")
         print(json.dumps(stats, ensure_ascii=False, indent=2))
         print(json.dumps(stats, ensure_ascii=False, indent=2))
         return 0
         return 0
     finally:
     finally: