""" System 7 Miner --------------- Module 1: Clean TikTok rows, apply lexicon-derived raw signals, and emit feature vectors. Core outputs: - tiktok10m_sys7_features.parquet - optional sys7_miner_metadata.json This script is designed to stream large CSVs in chunks (no 10M-row load), normalize text with ftfy, compute raw lexical scores from System 6 lexicons, and generate contextual embeddings via sentence-transformers. """ from __future__ import annotations import argparse import json import logging import math import re from collections import defaultdict from dataclasses import asdict, dataclass from pathlib import Path from typing import Dict, List, Optional, Sequence, Tuple import numpy as np import pandas as pd import pyarrow as pa import pyarrow.parquet as pq import pyarrow.dataset as ds try: from sentence_transformers import SentenceTransformer except ImportError: # pragma: no cover - handled at runtime SentenceTransformer = None # type: ignore try: import ftfy except ImportError: # pragma: no cover - handled at runtime ftfy = None # type: ignore try: from langdetect import detect as langdetect_detect except ImportError: # pragma: no cover - optional dependency langdetect_detect = None # type: ignore logger = logging.getLogger("sys7_miner") CONTROL_CHARS_RE = re.compile(r"[\u0000-\u0008\u000b\u000c\u000e-\u001f]") REPLACEMENT_RUN_RE = re.compile("�{2,}") WHITESPACE_RE = re.compile(r"\s+") HASHTAG_SPLIT_RE = re.compile(r"[A-Z]?[a-z]+|[0-9]+") TOKENIZER_RE = re.compile(r"[a-z0-9']+") @dataclass class MinerConfig: input_path: str output_parquet: str = "tiktok10m_sys7_features.parquet" lexicons: str = "system7_lexicons.json" phrase_lexicon: Optional[str] = None slang_lexicon: str = "slang_lexicon.json" label_orders: str = "label_orders.json" batch_size: int = 50_000 text_cap: int = 1024 embedding_model: str = "sentence-transformers/all-MiniLM-L6-v2" embedding_batch_size: int = 128 device: str = "cpu" workers: int = 2 metadata_path: Optional[str] = "sys7_miner_metadata.json" diagnostics_path: Optional[str] = "sys7_miner_diagnostics.json" min_chars: int = 10 min_tokens: int = 2 language_allowlist: Optional[List[str]] = None # e.g., ["en"] auto_language_detect: bool = False def load_json(path: str) -> Dict: with open(path, "r", encoding="utf-8") as f: return json.load(f) def auto_detect_language(text: str) -> Optional[str]: """Best-effort language detection using langdetect, if available.""" if not langdetect_detect: return None if not text: return None try: return str(langdetect_detect(text)) except Exception: return None def prepare_slang_map(raw) -> Dict[str, str]: """ Normalize slang lexicon to token -> canonical mapping. Accepts dict (already mapping) or list of objects with 'token' and optional 'canonical' keys. """ mapping: Dict[str, str] = {} if isinstance(raw, dict): for k, v in raw.items(): mapping[str(k).lower()] = str(v) if v is not None else str(k).lower() return mapping if isinstance(raw, list): for item in raw: if not isinstance(item, dict): continue token = item.get("token") or item.get("slang") or item.get("term") canonical = item.get("canonical") or item.get("value") or token if token: mapping[str(token).lower()] = str(canonical).lower() return mapping def orient_lexicons(raw_lexicons: Dict[str, Dict]) -> Dict[str, Dict[str, Dict[str, float]]]: """ Normalize lexicons into dim -> token -> {label: weight}. Accepts lexicons shaped as label -> token -> weight or weight objects. Ignores non-dict top-level entries (e.g., config/sources). """ def to_weight(val): if isinstance(val, dict) and "weight" in val: return val["weight"] return val oriented: Dict[str, Dict[str, Dict[str, float]]] = {} for dim, labels in raw_lexicons.items(): if not isinstance(labels, dict): continue # skip config/metadata blocks token_map: Dict[str, Dict[str, float]] = defaultdict(dict) for label, token_weights in labels.items(): if not isinstance(token_weights, dict): continue for token, weight in token_weights.items(): w = to_weight(weight) try: w_float = float(w) except (TypeError, ValueError): continue token_map[token.lower()][label] = w_float oriented[dim] = token_map return oriented def standardize_quotes(text: str) -> str: return ( text.replace("“", '"') .replace("”", '"') .replace("‘", "'") .replace("’", "'") .replace("´", "'") ) def split_hashtag(tag: str) -> List[str]: parts = HASHTAG_SPLIT_RE.findall(tag) if parts: return [p.lower() for p in parts] return [tag.lower()] def normalize_hashtags(raw: Optional[str], slang_map: Dict[str, str]) -> List[str]: if raw is None or (isinstance(raw, float) and math.isnan(raw)): return [] tokens: List[str] = [] if isinstance(raw, str): raw_parts = re.split(r"[,\s]+", raw) elif isinstance(raw, (list, tuple)): raw_parts = raw else: return tokens for part in raw_parts: if not part: continue cleaned = part.lstrip("#").strip() if not cleaned: continue for token in split_hashtag(cleaned): mapped = slang_map.get(token, token) if mapped: tokens.append(mapped) return tokens def fuse_text(description: Optional[str], hashtags: Optional[str], slang_map: Dict[str, str], cap: int) -> Tuple[str, List[str]]: desc = "" if description is None or (isinstance(description, float) and math.isnan(description)) else str(description) ht_tokens = normalize_hashtags(hashtags, slang_map) hashtag_text = " ".join(ht_tokens) fused = " ".join(filter(None, [desc, hashtag_text])).lower() if ftfy: fused = ftfy.fix_text(fused, normalization="NFC") fused = CONTROL_CHARS_RE.sub(" ", fused) fused = REPLACEMENT_RUN_RE.sub(" ", fused) fused = standardize_quotes(fused) fused = WHITESPACE_RE.sub(" ", fused).strip() if cap and len(fused) > cap: fused = fused[:cap] return fused, ht_tokens def tokenize(text: str, hashtag_tokens: Sequence[str]) -> List[str]: base_tokens = TOKENIZER_RE.findall(text) tokens = [t for t in base_tokens if t] tokens.extend([t.lower() for t in hashtag_tokens if t]) return tokens def compute_raw_scores(tokens: Sequence[str], lexicons: Dict[str, Dict[str, Dict[str, float]]], label_orders: Dict[str, List[str]]) -> Dict[str, List[float]]: scores: Dict[str, Dict[str, float]] = {dim: defaultdict(float) for dim in label_orders} for token in tokens: for dim, token_map in lexicons.items(): if token not in token_map: continue for label, weight in token_map[token].items(): scores[dim][label] += weight norm = max(1.0, math.log1p(len(tokens))) vectorized: Dict[str, List[float]] = {} for dim, order in label_orders.items(): vectorized[dim] = [scores[dim].get(label, 0.0) / norm for label in order] return vectorized def load_embedder(model_name: str, device: str) -> SentenceTransformer: if SentenceTransformer is None: raise ImportError("sentence-transformers is required. Install via `pip install sentence-transformers`.") logger.info("Loading embedding model %s on device=%s", model_name, device) return SentenceTransformer(model_name, device=device) def embed_texts(model: SentenceTransformer, texts: List[str], batch_size: int) -> List[List[float]]: embeddings = model.encode( texts, batch_size=batch_size, show_progress_bar=False, convert_to_numpy=True, ) return embeddings.tolist() def process_chunk( df: pd.DataFrame, lexicons: Dict[str, Dict[str, Dict[str, float]]], phrase_lex: Dict[str, Dict[str, float]], slang_map: Dict[str, str], label_orders: Dict[str, List[str]], embedder: SentenceTransformer, cfg: MinerConfig, ) -> Tuple[pa.Table, int, int]: records = [] texts: List[str] = [] filtered_lang = 0 filtered_short = 0 for row in df.itertuples(index=False): description = getattr(row, "description", None) if description is None and hasattr(row, "desc"): description = getattr(row, "desc", None) hashtags = getattr(row, "hashtags", None) lang = getattr(row, "language", None) or getattr(row, "lang", None) if cfg.auto_language_detect: # Only override missing/unknown language values lang_text_parts: List[str] = [] if description: lang_text_parts.append(str(description)) if hashtags: lang_text_parts.append(str(hashtags)) detected = auto_detect_language(" ".join(lang_text_parts)) if detected and (lang is None or str(lang).strip().lower() in ("", "und", "unknown")): lang = detected if cfg.language_allowlist and lang and str(lang).lower() not in cfg.language_allowlist: filtered_lang += 1 continue fused_text, ht_tokens = fuse_text(description, hashtags, slang_map, cfg.text_cap) tokens = tokenize(fused_text, ht_tokens) # Add matched phrases as tokens (if provided) if phrase_lex: txt_lower = fused_text.lower() for dim, phrases in phrase_lex.items(): for phrase in phrases.keys(): if phrase in txt_lower: tokens.append(phrase) if len(fused_text) < cfg.min_chars or len(tokens) < cfg.min_tokens: filtered_short += 1 continue texts.append(fused_text) raw_scores = compute_raw_scores(tokens, lexicons, label_orders) created_raw = getattr(row, "created_at", None) created_date = None if created_raw: try: created_int = int(created_raw) created_date = pd.to_datetime(created_int, unit="s", utc=True).date().isoformat() except Exception: try: created_date = str(pd.to_datetime(created_raw)).split(" ")[0] except Exception: created_date = None record = { "video_id": getattr(row, "video_id", None) or getattr(row, "aweme_id", None), "author_id": getattr(row, "author_id", None), "created_at": created_raw, "created_date": created_date, "language": lang, "clean_text": fused_text, "short_text": len(fused_text) < max(cfg.min_chars, 20), "tribe_scores_raw": raw_scores.get("tribe", []), "vibe_scores_raw": raw_scores.get("vibe", []), "commercial_scores_raw": raw_scores.get("commercial", []), "role_scores_raw": raw_scores.get("role", []), "format_scores_raw": raw_scores.get("format", []), "time_scores_raw": raw_scores.get("time", []), } records.append(record) if not records: return pa.Table.from_pylist([]), filtered_lang, filtered_short embeddings = embed_texts(embedder, texts, cfg.embedding_batch_size) for rec, emb in zip(records, embeddings): rec["embed_vec"] = emb table = pa.Table.from_pylist(records) return table, filtered_lang, filtered_short def write_metadata(cfg: MinerConfig, label_orders: Dict[str, List[str]], diagnostics: Optional[Dict[str, int]] = None) -> None: if not cfg.metadata_path: return meta = { "config": asdict(cfg), "label_orders": label_orders, "diagnostics": diagnostics or {}, } with open(cfg.metadata_path, "w", encoding="utf-8") as f: json.dump(meta, f, ensure_ascii=False, indent=2) logger.info("Wrote metadata to %s", cfg.metadata_path) def run(cfg: MinerConfig) -> None: logger.info("Starting sys7_miner with batch_size=%s", cfg.batch_size) # Prefer System 7 lexicons if configured; fall back to System 6 bundle if needed. lexicon_path = cfg.lexicons if not Path(lexicon_path).exists() and Path("system6_lexicons.json").exists(): lexicon_path = "system6_lexicons.json" logger.info("Lexicon bundle %s not found; falling back to %s", cfg.lexicons, lexicon_path) else: logger.info("Using lexicon bundle %s", lexicon_path) raw_lexicons = load_json(str(lexicon_path)) phrase_lex: Dict[str, Dict[str, float]] = {} if cfg.phrase_lexicon and Path(cfg.phrase_lexicon).exists(): try: phrase_bundle = load_json(cfg.phrase_lexicon) # Expecting {dim: {phrase: weight}} or wrapped under a "phrases" key if "phrases" in phrase_bundle: phrase_lex = {k: {p.lower(): float(w) for p, w in v.items()} for k, v in phrase_bundle["phrases"].items()} else: phrase_lex = {k: {p.lower(): float(w) for p, w in v.items()} for k, v in phrase_bundle.items()} logger.info("Loaded phrase lexicon from %s with dims=%s", cfg.phrase_lexicon, list(phrase_lex.keys())) except Exception as exc: # pragma: no cover - defensive logger.warning("Failed to load phrase lexicon %s: %s", cfg.phrase_lexicon, exc) slang_map = prepare_slang_map(load_json(cfg.slang_lexicon)) label_orders = load_json(cfg.label_orders) lexicons = orient_lexicons(raw_lexicons) embedder = load_embedder(cfg.embedding_model, cfg.device) writer: Optional[pq.ParquetWriter] = None total_rows = 0 kept_rows = 0 filtered_lang = 0 filtered_short = 0 path_obj = Path(cfg.input_path) is_parquet = path_obj.is_dir() or path_obj.suffix.lower() == ".parquet" def write_table(table: pa.Table): nonlocal writer, kept_rows if table.num_rows == 0: return if writer is None: writer = pq.ParquetWriter(cfg.output_parquet, table.schema, compression="snappy") writer.write_table(table) kept_rows += table.num_rows if is_parquet: dataset = ds.dataset(cfg.input_path, format="parquet") scanner = dataset.scanner(batch_size=cfg.batch_size) for idx, batch in enumerate(scanner.to_batches()): df = batch.to_pandas() total_rows += len(df) table, fl, fs = process_chunk(df, lexicons, phrase_lex, slang_map, label_orders, embedder, cfg) filtered_lang += fl filtered_short += fs write_table(table) logger.info("Processed parquet batch %s (input rows %s, kept total %s)", idx, len(df), kept_rows) else: usecols = ["video_id", "author_id", "created_at", "description", "hashtags", "language"] with pd.read_csv(cfg.input_path, chunksize=cfg.batch_size, usecols=lambda c: True, dtype=str, keep_default_na=False) as reader: for idx, chunk in enumerate(reader): total_rows += len(chunk) table, fl, fs = process_chunk(chunk, lexicons, phrase_lex, slang_map, label_orders, embedder, cfg) filtered_lang += fl filtered_short += fs write_table(table) logger.info("Processed CSV chunk %s (input rows %s, kept total %s)", idx, len(chunk), kept_rows) if writer: writer.close() diagnostics = { "input_rows": total_rows, "kept_rows": kept_rows, "filtered_language": filtered_lang, "filtered_short": filtered_short, } write_metadata(cfg, label_orders, diagnostics) if cfg.diagnostics_path: Path(cfg.diagnostics_path).write_text(json.dumps(diagnostics, indent=2), encoding="utf-8") logger.info( "Completed. Input rows: %s, kept: %s, filtered_language: %s, filtered_short: %s", total_rows, kept_rows, filtered_lang, filtered_short, ) def parse_args() -> MinerConfig: parser = argparse.ArgumentParser(description="System 7 Miner: clean + embed + score raw TikTok text.") parser.add_argument("--input-path", required=True, help="Path to CSV or Parquet directory (e.g., Data/hf10m/Small)") parser.add_argument("--output-parquet", default="tiktok10m_sys7_features.parquet", help="Output parquet path") parser.add_argument( "--lexicons", default=None, help="Preferred lexicon bundle (default: system7_lexicons.json if present, else system6_lexicons.json).", ) parser.add_argument( "--phrase-lexicon", default=None, help="Optional phrase lexicon JSON (e.g., mined phrases) to match substrings and include as tokens.", ) parser.add_argument( "--system6-lexicons", default="system6_lexicons.json", help="Fallback System 6 lexicon bundle (used if system7_lexicons.json is missing).", ) parser.add_argument("--slang-lexicon", default="slang_lexicon.json", help="Path to slang_lexicon.json") parser.add_argument("--label-orders", default="label_orders.json", help="Path to label_orders.json") parser.add_argument("--batch-size", type=int, default=50_000) parser.add_argument("--text-cap", type=int, default=1024) parser.add_argument("--embedding-model", default="sentence-transformers/all-MiniLM-L6-v2") parser.add_argument("--embedding-batch-size", type=int, default=128) parser.add_argument("--device", default="cpu", help="Embedding device: cpu or cuda") parser.add_argument("--workers", type=int, default=2, help="Reserved for future multiprocessing in tokenization") parser.add_argument("--metadata-path", default="sys7_miner_metadata.json") parser.add_argument("--diagnostics-path", default="sys7_miner_diagnostics.json") parser.add_argument("--min-chars", type=int, default=10, help="Minimum clean_text length to keep") parser.add_argument("--min-tokens", type=int, default=2, help="Minimum token count to keep") parser.add_argument( "--language-allowlist", default=None, help="Comma-separated list of language codes to keep (e.g., en,fr). " "Combined with optional --auto-language-detect to filter down to core languages.", ) parser.add_argument( "--auto-language-detect", action="store_true", help="If set, auto-detect language for rows with missing/unknown language using langdetect, when available.", ) args = parser.parse_args() lexicon_path = args.lexicons if not lexicon_path: # Default preference: system7_lexicons.json if present, else System 6 bundle. default_sys7 = Path("system7_lexicons.json") lexicon_path = str(default_sys7) if default_sys7.exists() else args.system6_lexicons return MinerConfig( input_path=args.input_path, output_parquet=args.output_parquet, lexicons=lexicon_path, phrase_lexicon=args.phrase_lexicon, slang_lexicon=args.slang_lexicon, label_orders=args.label_orders, batch_size=args.batch_size, text_cap=args.text_cap, embedding_model=args.embedding_model, embedding_batch_size=args.embedding_batch_size, device=args.device, workers=args.workers, metadata_path=args.metadata_path, diagnostics_path=args.diagnostics_path, min_chars=args.min_chars, min_tokens=args.min_tokens, language_allowlist=[s.strip().lower() for s in args.language_allowlist.split(",")] if args.language_allowlist else None, auto_language_detect=bool(args.auto_language_detect), ) def main() -> None: logging.basicConfig( level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s - %(message)s", ) cfg = parse_args() run(cfg) if __name__ == "__main__": main()