import logging import os import sys import sqlite3 import json import argparse import concurrent.futures from collections import defaultdict, deque from itertools import combinations from typing import Optional # Setup logging logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s') logger = logging.getLogger(__name__) # Get absolute paths BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) # Points to database/ directory PROJECT_ROOT = os.path.dirname(BASE_DIR) # Points to project root sys.path.insert(0, PROJECT_ROOT) # Add project root to Python path from database.paths import L2_DB, L3_DB, L3_SCHEMA, WEB_DB L2_DB_PATH = str(L2_DB) L3_DB_PATH = str(L3_DB) L3_BACKUP_PATH = f"{L3_DB_PATH}.bak" WEB_DB_PATH = str(WEB_DB) SCHEMA_PATH = str(L3_SCHEMA) def _get_existing_columns(conn, table_name): cur = conn.execute(f"PRAGMA table_info({table_name})") return {row[1] for row in cur.fetchall()} def _ensure_columns(conn, table_name, columns): existing = _get_existing_columns(conn, table_name) for col, col_type in columns.items(): if col in existing: continue conn.execute(f"ALTER TABLE {table_name} ADD COLUMN {col} {col_type}") def init_db(): """Initialize L3 database with new schema""" l3_dir = os.path.dirname(L3_DB_PATH) if not os.path.exists(l3_dir): os.makedirs(l3_dir) logger.info(f"Initializing L3 database at: {L3_DB_PATH}") conn = sqlite3.connect(L3_DB_PATH) try: with open(SCHEMA_PATH, 'r', encoding='utf-8') as f: schema_sql = f.read() conn.executescript(schema_sql) conn.commit() logger.info("✓ L3 schema created successfully") # Verify tables cursor = conn.cursor() cursor.execute("SELECT name FROM sqlite_master WHERE type='table' ORDER BY name") tables = [row[0] for row in cursor.fetchall()] logger.info(f"✓ Created {len(tables)} tables: {', '.join(tables)}") # Verify dm_player_features columns cursor.execute("PRAGMA table_info(dm_player_features)") columns = cursor.fetchall() logger.info(f"✓ dm_player_features has {len(columns)} columns") except Exception as e: logger.error(f"Error initializing L3 database: {e}") raise finally: conn.close() logger.info("L3 DB Initialized with new 5-tier architecture") def _get_team_players(): """Get list of steam_ids from Web App team lineups""" if not os.path.exists(WEB_DB_PATH): logger.warning(f"Web DB not found at {WEB_DB_PATH}, returning empty list") return set() try: conn = sqlite3.connect(WEB_DB_PATH) cursor = conn.cursor() columns = { row[1] for row in cursor.execute("PRAGMA table_info(team_lineups)") } if 'is_active' in columns: cursor.execute( """ SELECT player_ids_json FROM team_lineups WHERE is_active = 1 ORDER BY created_at DESC, id DESC LIMIT 1 """ ) else: cursor.execute( """ SELECT player_ids_json FROM team_lineups ORDER BY created_at DESC, id DESC LIMIT 1 """ ) rows = cursor.fetchall() steam_ids = set() for row in rows: if row[0]: try: ids = json.loads(row[0]) if isinstance(ids, list): steam_ids.update(ids) except json.JSONDecodeError: logger.warning(f"Failed to parse player_ids_json: {row[0]}") conn.close() logger.info(f"Found {len(steam_ids)} unique players in Team Lineups") return steam_ids except Exception as e: logger.error(f"Error reading Web DB: {e}") return set() def _get_match_date_range(steam_id: str, conn_l2: sqlite3.Connection): cursor = conn_l2.cursor() cursor.execute(""" SELECT MIN(m.start_time), MAX(m.start_time) FROM fact_match_players p JOIN fact_matches m ON p.match_id = m.match_id WHERE p.steam_id_64 = ? """, (steam_id,)) date_row = cursor.fetchone() first_match_date = date_row[0] if date_row and date_row[0] else None last_match_date = date_row[1] if date_row and date_row[1] else None return first_match_date, last_match_date def _build_player_record(steam_id: str): try: from database.L3.processors import ( BasicProcessor, TacticalProcessor, IntelligenceProcessor, MetaProcessor, CompositeProcessor ) conn_l2 = sqlite3.connect(L2_DB_PATH) conn_l2.row_factory = sqlite3.Row features = {} features.update(BasicProcessor.calculate(steam_id, conn_l2)) features.update(TacticalProcessor.calculate(steam_id, conn_l2)) features.update(IntelligenceProcessor.calculate(steam_id, conn_l2)) features.update(MetaProcessor.calculate(steam_id, conn_l2)) features.update(CompositeProcessor.calculate(steam_id, conn_l2, features)) match_count = _get_match_count(steam_id, conn_l2) round_count = _get_round_count(steam_id, conn_l2) first_match_date, last_match_date = _get_match_date_range(steam_id, conn_l2) conn_l2.close() return { "steam_id": steam_id, "features": features, "match_count": match_count, "round_count": round_count, "first_match_date": first_match_date, "last_match_date": last_match_date, "error": None, } except Exception as e: return { "steam_id": steam_id, "features": None, "match_count": 0, "round_count": 0, "first_match_date": None, "last_match_date": None, "error": str(e), } def _backup_l3_database(source_path=L3_DB_PATH, backup_path=L3_BACKUP_PATH): if not os.path.exists(source_path): return None source = sqlite3.connect(source_path) backup = sqlite3.connect(backup_path) try: source.backup(backup) result = backup.execute("PRAGMA quick_check").fetchone()[0] if result != 'ok': raise RuntimeError(f"L3 backup quick_check failed: {result}") finally: source.close() backup.close() logger.info("L3 backup created at %s", backup_path) return backup_path def main(force_all: bool = False, workers: int = 1, create_backup: bool = True): """ Main L3 feature building pipeline using modular processors """ logger.info("========================================") logger.info("Starting L3 Builder with 5-Tier Architecture") logger.info("========================================") if create_backup: _backup_l3_database() # 1. Ensure Schema is up to date init_db() # 2. Import processors try: from database.L3.processors import ( BasicProcessor, TacticalProcessor, IntelligenceProcessor, MetaProcessor, CompositeProcessor ) logger.info("✓ All 5 processors imported successfully") except ImportError as e: logger.error(f"Failed to import processors: {e}") return # 3. Connect to databases conn_l2 = sqlite3.connect(L2_DB_PATH) conn_l2.row_factory = sqlite3.Row conn_l3 = sqlite3.connect(L3_DB_PATH) try: conn_l3.execute("BEGIN IMMEDIATE") cursor_l2 = conn_l2.cursor() if force_all: logger.info("Force mode enabled: building L3 for all players in L2.") sql = """ SELECT DISTINCT steam_id_64 FROM dim_players ORDER BY steam_id_64 """ cursor_l2.execute(sql) else: team_players = _get_team_players() if not team_players: logger.warning("No players found in Team Lineups. Aborting L3 build.") return placeholders = ','.join(['?' for _ in team_players]) sql = f""" SELECT DISTINCT steam_id_64 FROM dim_players WHERE steam_id_64 IN ({placeholders}) ORDER BY steam_id_64 """ cursor_l2.execute(sql, list(team_players)) players = cursor_l2.fetchall() total_players = len(players) logger.info(f"Found {total_players} matching players in L2 to process") if total_players == 0: logger.warning("No matching players found in dim_players table") return success_count = 0 error_count = 0 processed_count = 0 if workers and workers > 1: steam_ids = [row[0] for row in players] with concurrent.futures.ProcessPoolExecutor(max_workers=workers) as executor: futures = [executor.submit(_build_player_record, sid) for sid in steam_ids] for future in concurrent.futures.as_completed(futures): result = future.result() processed_count += 1 if result.get("error"): error_count += 1 logger.error(f"Error processing player {result.get('steam_id')}: {result.get('error')}") else: _upsert_features( conn_l3, result["steam_id"], result["features"], result["match_count"], result["round_count"], None, result["first_match_date"], result["last_match_date"], ) success_count += 1 if processed_count % 2 == 0: logger.info(f"Progress: {processed_count}/{total_players} ({success_count} success, {error_count} errors)") else: for idx, row in enumerate(players, 1): steam_id = row[0] try: features = {} features.update(BasicProcessor.calculate(steam_id, conn_l2)) features.update(TacticalProcessor.calculate(steam_id, conn_l2)) features.update(IntelligenceProcessor.calculate(steam_id, conn_l2)) features.update(MetaProcessor.calculate(steam_id, conn_l2)) features.update(CompositeProcessor.calculate(steam_id, conn_l2, features)) match_count = _get_match_count(steam_id, conn_l2) round_count = _get_round_count(steam_id, conn_l2) first_match_date, last_match_date = _get_match_date_range(steam_id, conn_l2) _upsert_features(conn_l3, steam_id, features, match_count, round_count, conn_l2, first_match_date, last_match_date) success_count += 1 except Exception as e: error_count += 1 logger.error(f"Error processing player {steam_id}: {e}") if error_count <= 3: import traceback traceback.print_exc() continue processed_count = idx if processed_count % 2 == 0: logger.info(f"Progress: {processed_count}/{total_players} ({success_count} success, {error_count} errors)") if error_count: raise RuntimeError( f"L3 feature build failed for {error_count}/{total_players} players" ) processed_ids = [str(row[0]) for row in players] _update_percentiles(conn_l3, processed_ids) _rebuild_auxiliary_marts(conn_l2, conn_l3, processed_ids) _rebuild_team_marts(conn_l2, conn_l3, processed_ids) from database.L3.processors.narrative_processor import NarrativeProcessor narrative_counts = NarrativeProcessor.rebuild( conn_l2, conn_l3, processed_ids, ) logger.info("Narrative marts rebuilt: %s", narrative_counts) from database.L3.processors.discovery_processor import DiscoveryProcessor discovery_counts = DiscoveryProcessor.rebuild( conn_l2, conn_l3, processed_ids, ) logger.info("Discovery marts rebuilt: %s", discovery_counts) quick_check = conn_l3.execute("PRAGMA quick_check").fetchone()[0] if quick_check != 'ok': raise RuntimeError(f"L3 quick_check failed before commit: {quick_check}") conn_l3.commit() logger.info("========================================") logger.info(f"L3 Build Complete!") logger.info(f" Success: {success_count} players") logger.info(f" Errors: {error_count} players") logger.info(f" Total: {total_players} players") logger.info(f" Success Rate: {success_count/total_players*100:.1f}%") logger.info("========================================") except Exception as e: conn_l3.rollback() logger.error(f"Fatal error during L3 build: {e}") import traceback traceback.print_exc() raise finally: conn_l2.close() conn_l3.close() def _get_match_count(steam_id: str, conn_l2: sqlite3.Connection) -> int: """Get total match count for player""" cursor = conn_l2.cursor() cursor.execute(""" SELECT COUNT(*) FROM fact_match_players WHERE steam_id_64 = ? """, (steam_id,)) return cursor.fetchone()[0] def _get_round_count(steam_id: str, conn_l2: sqlite3.Connection) -> int: """Get total round count for player""" cursor = conn_l2.cursor() cursor.execute(""" SELECT COALESCE(SUM(round_total), 0) FROM fact_match_players WHERE steam_id_64 = ? """, (steam_id,)) return cursor.fetchone()[0] def _upsert_features(conn_l3: sqlite3.Connection, steam_id: str, features: dict, match_count: int, round_count: int, conn_l2: Optional[sqlite3.Connection], first_match_date=None, last_match_date=None): """ Insert or update player features in dm_player_features """ cursor_l3 = conn_l3.cursor() if first_match_date is None or last_match_date is None: if conn_l2 is not None: first_match_date, last_match_date = _get_match_date_range(steam_id, conn_l2) else: first_match_date = None last_match_date = None # Add metadata to features features['total_matches'] = match_count features['total_rounds'] = round_count features['first_match_date'] = first_match_date features['last_match_date'] = last_match_date # Build dynamic column list from features dict columns = ['steam_id_64'] + list(features.keys()) placeholders = ','.join(['?' for _ in columns]) columns_sql = ','.join(columns) # Build UPDATE SET clause for ON CONFLICT update_clauses = [f"{col}=excluded.{col}" for col in features.keys()] update_clause_sql = ','.join(update_clauses) values = [steam_id] + [features[k] for k in features.keys()] sql = f""" INSERT INTO dm_player_features ({columns_sql}) VALUES ({placeholders}) ON CONFLICT(steam_id_64) DO UPDATE SET {update_clause_sql}, last_updated=CURRENT_TIMESTAMP """ cursor_l3.execute(sql, values) def _rebuild_auxiliary_marts(conn_l2, conn_l3, steam_ids): """Rebuild player-grain marts used by profiles and trend APIs.""" if not steam_ids: return logger.info("Rebuilding L3 match, map and weapon marts") total_history = 0 total_maps = 0 total_weapons = 0 total_periods = 0 total_records = 0 for start in range(0, len(steam_ids), 400): chunk = steam_ids[start:start + 400] placeholders = ','.join('?' for _ in chunk) for table in ( 'dm_player_match_history', 'dm_player_map_stats', 'dm_player_weapon_stats', 'dm_player_period_stats', 'dm_player_records', ): conn_l3.execute( f"DELETE FROM {table} WHERE steam_id_64 IN ({placeholders})", chunk, ) history_rows = conn_l2.execute( f""" SELECT mp.steam_id_64, mp.match_id, m.start_time, mp.rating, mp.kd_ratio, mp.adr, mp.kast, mp.is_win, m.map_name, mp.kills, mp.deaths, mp.headshot_count, ( SELECT AVG(teammate.rating) FROM fact_match_players teammate WHERE teammate.match_id = mp.match_id AND teammate.team_id = mp.team_id AND teammate.steam_id_64 != mp.steam_id_64 ) AS teammate_avg_rating FROM fact_match_players mp JOIN fact_matches m ON m.match_id = mp.match_id WHERE mp.steam_id_64 IN ({placeholders}) ORDER BY mp.steam_id_64, m.start_time, mp.match_id """, chunk, ).fetchall() history_values = [] player_state = defaultdict(lambda: { 'sequence': 0, 'rating_sum': 0.0, 'recent': deque(maxlen=10), }) for row in history_rows: steam_id = str(row[0]) state = player_state[steam_id] rating = float(row[3] or 0.0) state['sequence'] += 1 state['rating_sum'] += rating state['recent'].append(rating) history_values.append(( steam_id, row[1], row[2], state['sequence'], row[3], row[4], row[5], row[6], row[7], row[8], None, row[12], state['rating_sum'] / state['sequence'], sum(state['recent']) / len(state['recent']), )) conn_l3.executemany( """ INSERT INTO dm_player_match_history ( steam_id_64, match_id, match_date, match_sequence, rating, kd_ratio, adr, kast, is_win, map_name, opponent_avg_elo, teammate_avg_rating, cumulative_rating, rolling_10_rating ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, history_values, ) total_history += len(history_values) map_rows = conn_l2.execute( f""" SELECT mp.steam_id_64, m.map_name, COUNT(*) AS matches, SUM(CASE WHEN mp.is_win = 1 THEN 1 ELSE 0 END) AS wins, AVG(mp.rating) AS avg_rating, AVG(mp.kd_ratio) AS avg_kd, AVG(mp.adr) AS avg_adr, AVG(mp.kast) AS avg_kast, MAX(mp.rating) AS best_rating, MIN(mp.rating) AS worst_rating FROM fact_match_players mp JOIN fact_matches m ON m.match_id = mp.match_id WHERE mp.steam_id_64 IN ({placeholders}) AND m.map_name IS NOT NULL AND m.map_name != '' GROUP BY mp.steam_id_64, m.map_name """, chunk, ).fetchall() map_values = [ tuple(row[:4]) + ( (row[3] or 0) / row[2] if row[2] else 0.0, ) + tuple(row[4:]) for row in map_rows ] conn_l3.executemany( """ INSERT INTO dm_player_map_stats ( steam_id_64, map_name, matches, wins, win_rate, avg_rating, avg_kd, avg_adr, avg_kast, best_rating, worst_rating ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, map_values, ) total_maps += len(map_values) round_counts = { str(row[0]): int(row[1] or 0) for row in conn_l2.execute( f""" SELECT steam_id_64, SUM(round_total) FROM fact_match_players WHERE steam_id_64 IN ({placeholders}) GROUP BY steam_id_64 """, chunk, ) } weapon_rows = conn_l2.execute( f""" SELECT attacker_steam_id, weapon, COUNT(*) AS total_kills, SUM(CASE WHEN is_headshot = 1 THEN 1 ELSE 0 END) AS total_headshots, COUNT(DISTINCT match_id || ':' || round_num) AS usage_rounds FROM fact_round_events WHERE event_type = 'kill' AND attacker_steam_id IN ({placeholders}) AND weapon IS NOT NULL AND weapon != '' GROUP BY attacker_steam_id, weapon """, chunk, ).fetchall() weapon_values = [] for row in weapon_rows: rounds = round_counts.get(str(row[0]), 0) kills = int(row[2] or 0) headshots = int(row[3] or 0) usage_rounds = int(row[4] or 0) hs_rate = headshots / kills if kills else 0.0 usage_rate = usage_rounds / rounds if rounds else 0.0 kills_per_round = kills / rounds if rounds else 0.0 effectiveness = kills / usage_rounds if usage_rounds else 0.0 weapon_values.append(( str(row[0]), row[1], kills, headshots, hs_rate, usage_rounds, usage_rate, kills_per_round, effectiveness, )) conn_l3.executemany( """ INSERT INTO dm_player_weapon_stats ( steam_id_64, weapon_name, total_kills, total_headshots, hs_rate, usage_rounds, usage_rate, avg_kills_per_round, effectiveness_score ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) """, weapon_values, ) total_weapons += len(weapon_values) period_values = _calculate_period_rows(history_rows) conn_l3.executemany( """ INSERT INTO dm_player_period_stats ( steam_id_64, period_key, period_label, period_start, period_end, matches, wins, win_rate, avg_rating, avg_kd, avg_adr, avg_kast, total_kills, total_deaths, sample_reliable ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, period_values, ) total_periods += len(period_values) record_values = _calculate_record_rows(history_rows) conn_l3.executemany( """ INSERT INTO dm_player_records ( steam_id_64, record_key, record_label, record_value, match_id, map_name, match_date ) VALUES (?, ?, ?, ?, ?, ?, ?) """, record_values, ) total_records += len(record_values) logger.info( "Auxiliary marts rebuilt: %s history, %s map, %s weapon, " "%s period, %s record rows", total_history, total_maps, total_weapons, total_periods, total_records, ) def _group_player_match_rows(history_rows): grouped = defaultdict(list) for row in history_rows: grouped[str(row['steam_id_64'])].append(row) for rows in grouped.values(): rows.sort(key=lambda row: (row['start_time'] or 0, row['match_id'])) return grouped def _safe_average(rows, key): values = [float(row[key]) for row in rows if row[key] is not None] return sum(values) / len(values) if values else None def _calculate_period_rows(history_rows): result = [] for steam_id, all_rows in _group_player_match_rows(history_rows).items(): latest_time = max(int(row['start_time'] or 0) for row in all_rows) period_groups = [ ('career', '生涯', all_rows), ('last_10', '最近 10 场', all_rows[-10:]), ('last_20', '最近 20 场', all_rows[-20:]), ('last_30', '最近 30 场', all_rows[-30:]), ( 'days_30', '最近 30 天', [ row for row in all_rows if int(row['start_time'] or 0) >= latest_time - 30 * 86400 ], ), ( 'days_90', '最近 90 天', [ row for row in all_rows if int(row['start_time'] or 0) >= latest_time - 90 * 86400 ], ), ] for period_key, period_label, rows in period_groups: if not rows: continue matches = len(rows) wins = sum(1 for row in rows if row['is_win']) kills = sum(int(row['kills'] or 0) for row in rows) deaths = sum(int(row['deaths'] or 0) for row in rows) result.append(( steam_id, period_key, period_label, min(int(row['start_time'] or 0) for row in rows), max(int(row['start_time'] or 0) for row in rows), matches, wins, wins / matches, _safe_average(rows, 'rating'), kills / deaths if deaths else float(kills), _safe_average(rows, 'adr'), _safe_average(rows, 'kast'), kills, deaths, 1 if matches >= 10 else 0, )) return result def _calculate_record_rows(history_rows): result = [] metric_definitions = ( ('highest_rating', '最高 Rating', 'rating'), ('most_kills', '单场最多击杀', 'kills'), ('highest_adr', '单场最高 ADR', 'adr'), ('highest_kd', '单场最高 K/D', 'kd_ratio'), ('most_headshots', '单场最多爆头', 'headshot_count'), ) for steam_id, rows in _group_player_match_rows(history_rows).items(): for record_key, record_label, field in metric_definitions: candidates = [row for row in rows if row[field] is not None] if not candidates: continue best = max( candidates, key=lambda row: ( float(row[field]), int(row['start_time'] or 0), ), ) result.append(( steam_id, record_key, record_label, float(best[field]), best['match_id'], best['map_name'], best['start_time'], )) longest_streak = 0 current_streak = 0 streak_end = None for row in rows: if row['is_win']: current_streak += 1 if current_streak >= longest_streak: longest_streak = current_streak streak_end = row else: current_streak = 0 if streak_end is not None: result.append(( steam_id, 'longest_win_streak', '最长连胜', float(longest_streak), streak_end['match_id'], streak_end['map_name'], streak_end['start_time'], )) return result def _rebuild_team_marts(conn_l2, conn_l3, steam_ids): if not steam_ids: return conn_l3.execute('DELETE FROM dm_duo_stats') conn_l3.execute('DELETE FROM dm_lineup_stats') all_rows = [] for start in range(0, len(steam_ids), 400): chunk = steam_ids[start:start + 400] placeholders = ','.join('?' for _ in chunk) all_rows.extend(conn_l2.execute( f""" SELECT p.match_id, p.steam_id_64, CASE WHEN p.group_id IN (1, 2) THEN p.group_id WHEN p.team_id IN (1, 2) THEN p.team_id END AS team_key, p.rating, p.is_win, m.start_time FROM fact_match_players p JOIN fact_matches m ON m.match_id = p.match_id WHERE p.steam_id_64 IN ({placeholders}) """, chunk, ).fetchall()) match_teams = defaultdict(list) for row in all_rows: if row['team_key'] is None: continue match_teams[(row['match_id'], row['team_key'])].append(row) duo_accumulator = {} lineup_accumulator = {} for rows in match_teams.values(): players = { str(row['steam_id_64']): row for row in rows } player_ids = sorted(players) if len(player_ids) < 2: continue is_win = bool(next(iter(players.values()))['is_win']) match_date = int(next(iter(players.values()))['start_time'] or 0) for player_a, player_b in combinations(player_ids, 2): key = (player_a, player_b) accumulator = duo_accumulator.setdefault(key, { 'matches': 0, 'wins': 0, 'rating_a': 0.0, 'rating_b': 0.0, 'first': match_date, 'last': match_date, }) accumulator['matches'] += 1 accumulator['wins'] += int(is_win) accumulator['rating_a'] += float(players[player_a]['rating'] or 0) accumulator['rating_b'] += float(players[player_b]['rating'] or 0) accumulator['first'] = min(accumulator['first'], match_date) accumulator['last'] = max(accumulator['last'], match_date) for size in range(2, min(5, len(player_ids)) + 1): for selected_ids in combinations(player_ids, size): lineup_key = '|'.join(selected_ids) accumulator = lineup_accumulator.setdefault(lineup_key, { 'player_ids': selected_ids, 'matches': 0, 'wins': 0, 'rating': 0.0, 'first': match_date, 'last': match_date, }) accumulator['matches'] += 1 accumulator['wins'] += int(is_win) accumulator['rating'] += ( sum(float(players[steam_id]['rating'] or 0) for steam_id in selected_ids) / len(selected_ids) ) accumulator['first'] = min(accumulator['first'], match_date) accumulator['last'] = max(accumulator['last'], match_date) duo_values = [] for (player_a, player_b), value in duo_accumulator.items(): matches = value['matches'] avg_a = value['rating_a'] / matches avg_b = value['rating_b'] / matches duo_values.append(( player_a, player_b, matches, value['wins'], value['wins'] / matches, avg_a, avg_b, (avg_a + avg_b) / 2, value['first'], value['last'], 1 if matches >= 5 else 0, )) conn_l3.executemany( """ INSERT INTO dm_duo_stats ( steam_id_a, steam_id_b, matches, wins, win_rate, avg_rating_a, avg_rating_b, avg_combined_rating, first_match_date, last_match_date, sample_reliable ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, duo_values, ) lineup_values = [] for lineup_key, value in lineup_accumulator.items(): matches = value['matches'] lineup_values.append(( lineup_key, json.dumps(value['player_ids']), len(value['player_ids']), matches, value['wins'], value['wins'] / matches, value['rating'] / matches, value['first'], value['last'], 1 if matches >= 3 else 0, )) conn_l3.executemany( """ INSERT INTO dm_lineup_stats ( lineup_key, player_ids_json, player_count, matches, wins, win_rate, avg_team_rating, first_match_date, last_match_date, sample_reliable ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, lineup_values, ) logger.info( 'Team marts rebuilt: %s duos, %s lineups', len(duo_values), len(lineup_values), ) def _update_percentiles(conn_l3, steam_ids): """Calculate a real percentile among eligible players in this build.""" if not steam_ids: return score_rows = [] for start in range(0, len(steam_ids), 400): chunk = steam_ids[start:start + 400] placeholders = ','.join('?' for _ in chunk) conn_l3.execute( f""" UPDATE dm_player_features SET tier_percentile = NULL WHERE steam_id_64 IN ({placeholders}) """, chunk, ) score_rows.extend(conn_l3.execute( f""" SELECT steam_id_64, score_overall FROM dm_player_features WHERE steam_id_64 IN ({placeholders}) AND score_overall > 0 """, chunk, ).fetchall()) if not score_rows: return scores = [float(row[1]) for row in score_rows] percentile_values = [] for row in score_rows: score = float(row[1]) percentile = sum(value <= score for value in scores) / len(scores) * 100 percentile_values.append((round(percentile, 2), str(row[0]))) conn_l3.executemany( """ UPDATE dm_player_features SET tier_percentile = ? WHERE steam_id_64 = ? """, percentile_values, ) logger.info("Updated percentiles for %s eligible players", len(score_rows)) def _parse_args(): parser = argparse.ArgumentParser() parser.add_argument("--force", action="store_true") parser.add_argument("--workers", type=int, default=1) parser.add_argument("--no-backup", action="store_true") return parser.parse_args() if __name__ == "__main__": args = _parse_args() main( force_all=args.force, workers=args.workers, create_backup=not args.no_backup, )