yumeng 1 hari lalu
induk
melakukan
144851c268
1 mengubah file dengan 9 tambahan dan 2 penghapusan
  1. 9 2
      scripts/honor_backfill_oaid.py

+ 9 - 2
scripts/honor_backfill_oaid.py

@@ -180,7 +180,7 @@ class RespRedisClient:
         self._read_response()
 
 
-def fetch_rows(conn, trace_ids: List[str]) -> Dict[str, dict]:
+def fetch_rows(conn, trace_ids: List[str], progress_every: int) -> Dict[str, dict]:
     rows: Dict[str, dict] = {}
     chunk_size = 500
     sql = """
@@ -188,10 +188,16 @@ def fetch_rows(conn, trace_ids: List[str]) -> Dict[str, dict]:
         FROM honor_ad_bid_events
         WHERE media_trace_id IN ({})
     """
+    total_chunks = (len(trace_ids) + chunk_size - 1) // chunk_size
     with conn.cursor() as cur:
         for i in range(0, len(trace_ids), chunk_size):
             chunk = trace_ids[i:i + chunk_size]
             placeholders = ",".join(["%s"] * len(chunk))
+            chunk_no = i // chunk_size + 1
+            if chunk_no == 1 or chunk_no == total_chunks or (progress_every > 0 and chunk_no % max(progress_every // chunk_size, 1) == 0):
+                print(
+                    f"[db] fetching chunk={chunk_no}/{total_chunks} chunkSize={len(chunk)} rowsLoaded={len(rows)}"
+                )
             cur.execute(sql.format(placeholders), chunk)
             for qk, media_trace_id, media_params in cur.fetchall():
                 rows[media_trace_id] = {
@@ -199,6 +205,7 @@ def fetch_rows(conn, trace_ids: List[str]) -> Dict[str, dict]:
                     "trace_id": media_trace_id,
                     "media_params": json.loads(media_params) if media_params else {},
                 }
+    print(f"[db] fetched rows={len(rows)}")
     return rows
 
 
@@ -270,7 +277,7 @@ def main() -> int:
     )
 
     try:
-        db_rows = fetch_rows(conn, trace_ids)
+        db_rows = fetch_rows(conn, trace_ids, args.progress_every)
         stats = {
             "log_mappings": len(items),
             "db_found": 0,