Miner / sys7_miner.py
Imaginethat's picture
Upload sys7_miner.py
577cc97 verified
Raw History Blame Contribute Delete
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']+")
@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()