reranker / src /preprocess /load.py
Hemprasad Badgujar
Add bm25 prefilter stage; unify candidates cache; validate reasoning quality
cd5f9d4
Raw History Blame Contribute Delete
13.5 kB
"""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