Spaces:
Paused
Paused
Download sys7_miner.py from Imaginethat/Miner: direct link, hf CLI and curl.
- Browser
- Download file 20.3 kB
-
https://huggingface.co/spaces/Imaginethat/Miner/resolve/main/sys7_miner.py
- Command line
-
hf download hf://spaces/Imaginethat/Miner/sys7_miner.py
-
curl -L -o sys7_miner.py https://huggingface.co/spaces/Imaginethat/Miner/resolve/main/sys7_miner.py
20.3 kB
| """ | |
| 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']+") | |
| 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() | |