"""Step 02 sub-phase 1 — load candidates.jsonl into one flat ranking frame. Each JSON line becomes one row. ``_flatten`` hoists every scalar (profile + the 23 redrob signals) to a top-level column and keeps the five nested arrays (career_history / education / skills / certifications / languages) structured for later enrichment. WHY one place covers every field: a key missed here is invisibly absent in every downstream stage — so ``_flatten`` is the field-coverage contract and is regression-tested. No row is ever dropped on content; only blank or undecodable lines are skipped. For the real, unlimited load of the default ``candidates.jsonl``, the result is ALSO cached as Parquet at ``LOADED_CANDIDATES_CACHE`` (see ``load_candidates``'s docstring) — shared by every caller (CLI and Streamlit alike), since it lives at this layer rather than one layer up in ``run_preprocess``. """ from __future__ import annotations import gzip from pathlib import Path from typing import IO, Any import orjson import polars as pl from common.io import atomic_write_parquet from common.logging import get_logger, step from common.paths import DEFAULT_CANDIDATES, LOADED_CANDIDATES_CACHE # The dataset's documented sentinel for "no signal" (github / offer history). Preserved verbatim on load; # it is cleaned in the later materialize stage, never here. UNKNOWN_SENTINEL = -1.0 # Explicit nested dtypes so every list column is typed even when empty (a row/pool with no skills must # still be List(Struct), not List(Null) — otherwise downstream struct-field access breaks). This makes the # loaded schema deterministic instead of inference-dependent. _STR_LIST = pl.List(pl.String) _SCHEMA_OVERRIDES: dict[str, pl.DataType] = { "career_history": pl.List(pl.Struct({ "company": pl.String, "title": pl.String, "start_date": pl.String, "end_date": pl.String, "duration_months": pl.Int64, "is_current": pl.Boolean, "industry": pl.String, "company_size": pl.String, "description": pl.String, })), "education": pl.List(pl.Struct({ "institution": pl.String, "degree": pl.String, "field_of_study": pl.String, "start_year": pl.Int64, "end_year": pl.Int64, "grade": pl.String, "tier": pl.String, })), "skills": pl.List(pl.Struct({ "name": pl.String, "proficiency": pl.String, "endorsements": pl.Int64, "duration_months": pl.Int64, })), "certifications": pl.List(pl.Struct({"name": pl.String, "issuer": pl.String, "year": pl.Int64})), "languages": pl.List(pl.Struct({"language": pl.String, "proficiency": pl.String})), "skill_assessment_scores": pl.List(pl.Struct({"name": pl.String, "score": pl.Float64})), "flat_skills": _STR_LIST, "flat_languages": _STR_LIST, "flat_certifications": _STR_LIST, "flat_job_titles": _STR_LIST, "flat_career_industries": _STR_LIST, "flat_company_sizes": _STR_LIST, "flat_institutions": _STR_LIST, "flat_institution_tiers": _STR_LIST, "flat_degrees": _STR_LIST, "verified_skills": _STR_LIST, } def _as_float(value: Any, default: float = 0.0) -> float: """Coerce to float, tolerating None/blank; a real -1 sentinel passes straight through.""" if value is None or value == "": return default try: return float(value) except (TypeError, ValueError): return default def _as_int(value: Any, default: int = 0) -> int: """Coerce to int, tolerating None/blank/float-encoded ints.""" if value is None or value == "": return default try: return int(value) except (TypeError, ValueError): return default def to_date(expr: pl.Expr) -> pl.Expr: """Parse a ``YYYY-MM-DD`` column to Date (non-strict — bad/empty values become null). Shared date-parsing helper lives with the data/load module (the leaf) so the validate and later feature stages can use it without an import cycle through the pipeline orchestrator. """ return expr.cast(pl.String, strict=False).str.to_date("%Y-%m-%d", strict=False) def _flatten(record: dict[str, Any]) -> dict[str, Any]: """Map one raw candidate record to a flat row covering EVERY schema field. Scalars from ``profile`` and ``redrob_signals`` are hoisted to top-level columns; the salary range is un-nested to min/max; the dynamic ``skill_assessment_scores`` map is normalized to a typed list[{name, score}] (so it stays columnar instead of a free-form dict); the five document arrays are kept structured AND projected to derived name-only ``flat_*`` lists (cheap, filter/feature-friendly columns that mirror one another). All derived here, at the single decode point, never recomputed. """ profile = record.get("profile") or {} signals = record.get("redrob_signals") or {} salary = signals.get("expected_salary_range_inr_lpa") or {} career_history = record.get("career_history") or [] education = record.get("education") or [] skills = record.get("skills") or [] certifications = record.get("certifications") or [] languages = record.get("languages") or [] # Normalize skill->score map to a list of structs — dynamic keys can't be a stable columnar type. raw_assessments = signals.get("skill_assessment_scores") or {} skill_assessment_scores = [ {"name": name, "score": _as_float(score)} for name, score in raw_assessments.items() ] # The skill names that carry an assessment-score entry — a "verified by platform test" quality proxy # used by validate/materialize. Derived here (at the single decode point), not recomputed downstream. verified_skills = list(raw_assessments.keys()) return { "candidate_id": record.get("candidate_id"), # --- profile (hoisted scalars) --- "anonymized_name": profile.get("anonymized_name", ""), "headline": profile.get("headline", ""), "summary": profile.get("summary", ""), "location": profile.get("location", ""), "country": profile.get("country", ""), "years_of_experience": _as_float(profile.get("years_of_experience")), "current_title": profile.get("current_title", ""), "current_company": profile.get("current_company", ""), "current_company_size": profile.get("current_company_size", ""), "current_industry": profile.get("current_industry", ""), # --- nested document arrays (kept structured) --- "career_history": career_history, "education": education, "skills": skills, "certifications": certifications, "languages": languages, # --- derived name-only list projections of the nested arrays (filter/feature-friendly) --- "flat_skills": [(skill.get("name") or "") for skill in skills], "flat_languages": [(lang.get("language") or "") for lang in languages], "flat_certifications": [(cert.get("name") or "") for cert in certifications], "flat_job_titles": [(role.get("title") or "") for role in career_history], "flat_career_industries": [(role.get("industry") or "") for role in career_history], "flat_company_sizes": [(role.get("company_size") or "") for role in career_history], "flat_institutions": [(degree.get("institution") or "") for degree in education], "flat_institution_tiers": [(degree.get("tier") or "") for degree in education], "flat_degrees": [(degree.get("degree") or "") for degree in education], # --- redrob signals (hoisted scalars) --- "profile_completeness_score": _as_float(signals.get("profile_completeness_score")), "signup_date": signals.get("signup_date", ""), "last_active_date": signals.get("last_active_date", ""), "open_to_work_flag": bool(signals.get("open_to_work_flag", False)), "profile_views_received_30d": _as_int(signals.get("profile_views_received_30d")), "applications_submitted_30d": _as_int(signals.get("applications_submitted_30d")), "recruiter_response_rate": _as_float(signals.get("recruiter_response_rate")), "avg_response_time_hours": _as_float(signals.get("avg_response_time_hours")), "skill_assessment_scores": skill_assessment_scores, "verified_skills": verified_skills, "connection_count": _as_int(signals.get("connection_count")), "endorsements_received": _as_int(signals.get("endorsements_received")), "notice_period_days": _as_int(signals.get("notice_period_days")), "expected_salary_min": _as_float(salary.get("min")), "expected_salary_max": _as_float(salary.get("max")), "preferred_work_mode": signals.get("preferred_work_mode", ""), "willing_to_relocate": bool(signals.get("willing_to_relocate", False)), # github / offer keep their -1 "unknown" sentinel verbatim (cleaned later, not here): "github_activity_score": _as_float(signals.get("github_activity_score"), UNKNOWN_SENTINEL), "search_appearance_30d": _as_int(signals.get("search_appearance_30d")), "saved_by_recruiters_30d": _as_int(signals.get("saved_by_recruiters_30d")), "interview_completion_rate": _as_float(signals.get("interview_completion_rate")), "offer_acceptance_rate": _as_float(signals.get("offer_acceptance_rate"), UNKNOWN_SENTINEL), "verified_email": bool(signals.get("verified_email", False)), "verified_phone": bool(signals.get("verified_phone", False)), "linkedin_connected": bool(signals.get("linkedin_connected", False)), } def _open_text(path: Path) -> IO[str]: """Open a .jsonl or .jsonl.gz transparently as a UTF-8 text stream.""" if path.suffix == ".gz": return gzip.open(path, "rt", encoding="utf-8") return path.open("r", encoding="utf-8") def _cache_is_fresh(cache_path: Path, source_path: Path) -> bool: """The Parquet cache is usable only if it exists and is at least as new as the source JSONL.""" return cache_path.exists() and source_path.exists() and cache_path.stat().st_mtime >= source_path.stat().st_mtime def load_candidates(path: str | Path, limit: int | None = None) -> pl.DataFrame: """Read candidates.jsonl → a flat Polars frame, one fully-covered row per candidate. Args: path: Path to ``candidates.jsonl`` (or ``.jsonl.gz``). limit: If set, stop after this many decoded records (for fast smoke runs). Returns: A Polars DataFrame; ``candidate_id`` is non-null on every row. Blank and undecodable lines are skipped (counted in the step metrics), never raised — content is never dropped. For the real, unlimited load of the DEFAULT candidates path, the result is cached as Parquet at ``LOADED_CANDIDATES_CACHE`` — written once by ``rank.py --setup`` (bakes into the Docker image) or on first use, then read directly on every later call (~13-20s → ~1-2s) instead of re-parsing the 465MB JSONL. Every caller shares this cache (this function sits BELOW ``run_preprocess``, so both the CLI and Streamlit's per-phase explorer — which calls ``load_candidates`` directly — benefit). Not used for a custom ``path`` or a ``limit``'d dev/smoke read — those aren't the cached full dataset. """ source = Path(path) use_cache = limit is None and source.resolve() == DEFAULT_CANDIDATES.resolve() log = get_logger("preprocess") if use_cache and _cache_is_fresh(LOADED_CANDIDATES_CACHE, DEFAULT_CANDIDATES): try: frame = pl.read_parquet(LOADED_CANDIDATES_CACHE) log.info("preprocess.load_cache_hit", path=str(LOADED_CANDIDATES_CACHE), rows=frame.height) return frame except Exception as failure: # a corrupt/partial cache must never break the run — fall through to a live parse log.warning("preprocess.load_cache_unreadable", path=str(LOADED_CANDIDATES_CACHE), error=str(failure)) rows: list[dict[str, Any]] = [] blank_lines = 0 malformed_lines = 0 with step("preprocess.load", source=source.name) as metrics: with _open_text(source) as stream: for raw_line in stream: stripped = raw_line.strip() if not stripped: blank_lines += 1 # skip blank separators continue try: record = orjson.loads(stripped) except orjson.JSONDecodeError: malformed_lines += 1 # skip malformed line, do not abort the whole load continue rows.append(_flatten(record)) if limit is not None and len(rows) >= limit: break # Explicit nested dtypes (always typed, even when empty) + infer the rest. infer_schema_length=None # scans ALL rows so any non-overridden column resolves from real data, not just the first row. frame = pl.DataFrame(rows, schema_overrides=_SCHEMA_OVERRIDES, infer_schema_length=None) metrics["rows_out"] = frame.height metrics["blank"] = blank_lines metrics["malformed"] = malformed_lines if use_cache: try: atomic_write_parquet(frame, LOADED_CANDIDATES_CACHE) log.info("preprocess.load_cache_written", path=str(LOADED_CANDIDATES_CACHE)) except Exception as failure: # a cache-write failure must never break the run log.warning("preprocess.load_cache_write_failed", path=str(LOADED_CANDIDATES_CACHE), error=str(failure)) return frame