Hemprasad Badgujar
Add bm25 prefilter stage; unify candidates cache; validate reasoning quality
cd5f9d4 Download src/preprocess/load.py from hembad/reranker: direct link, hf CLI and curl.
- Browser
- Download file 13.5 kB
-
https://huggingface.co/spaces/hembad/reranker/resolve/main/src/preprocess/load.py
- Command line
-
hf download hf://spaces/hembad/reranker/src/preprocess/load.py
-
curl -L -o load.py https://huggingface.co/spaces/hembad/reranker/resolve/main/src/preprocess/load.py
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 | |