File size: 13,485 Bytes
28c12aa cd5f9d4 28c12aa 7d4eb85 28c12aa cd5f9d4 28c12aa 700e896 28c12aa 700e896 28c12aa 700e896 28c12aa 700e896 28c12aa 700e896 28c12aa 700e896 28c12aa 700e896 28c12aa 700e896 28c12aa cd5f9d4 28c12aa cd5f9d4 28c12aa cd5f9d4 28c12aa 7d4eb85 28c12aa 700e896 28c12aa cd5f9d4 28c12aa | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 | """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
|