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