reranker / src /preprocess /materialize.py
Hemprasad Badgujar
Add CPU threading, T1-only LLM polish, experiment script
a41b023
Raw History Blame Contribute Delete
10.8 kB
"""Preprocess sub-phase — materialize: one vectorized pass to features + the embedded text body.
Lean scope (this build): structural counts (from the raw ``flat_skills`` etc.), ``experience_band``
computed inline from years-of-experience, sentinel-clean platform signals, the salary band, and the
deterministic templated ``text`` that retrieval embeds (identity + experience + skills + education — NO
behavioral/platform signals) plus its content hash. The 7 quality pillars + ``candidate_quality`` and the
rerank-feature legs are deferred to their own later stages (they are ranking-side features).
"""
from __future__ import annotations
import hashlib
from datetime import date
import polars as pl
from common.logging import step
from .load import to_date
# Proficiency levels that count as a "core" (strong) skill.
_PROFICIENT = ["advanced", "expert"]
# Notice period (days) under which a candidate counts as "available soon".
_AVAILABLE_SOON_DAYS = 30
# Float inputs left at source precision — every other derived float is rounded to 3 dp in one pass.
_RAW_FLOAT_PASSTHROUGH = {
"years_of_experience", "profile_completeness_score", "recruiter_response_rate",
"avg_response_time_hours", "expected_salary_min", "expected_salary_max",
"github_activity_score", "interview_completion_rate", "offer_acceptance_rate",
}
def _text(column: str) -> pl.Expr:
"""Null-safe string view of a column (null → empty)."""
return pl.col(column).cast(pl.String).fill_null("")
def _experience_band_expr() -> pl.Expr:
"""Left-closed seniority band from years-of-experience (computed inline, no enrich)."""
yoe = pl.col("years_of_experience")
return (
pl.when(yoe < 3).then(pl.lit("junior"))
.when(yoe < 6).then(pl.lit("mid"))
.when(yoe < 9).then(pl.lit("senior"))
.when(yoe < 13).then(pl.lit("lead"))
.otherwise(pl.lit("principal"))
.alias("experience_band")
)
def _count_exprs() -> list[pl.Expr]:
"""Structural counts over the nested arrays + raw flat lists (skills via raw ``flat_skills``)."""
field = pl.element().struct.field
durations = pl.col("career_history").list.eval(field("duration_months").cast(pl.Int64).fill_null(0))
assessed = pl.col("skill_assessment_scores").list.eval(
pl.when(field("score") >= 0).then(field("score")).otherwise(None)
)
advanced = pl.col("skills").list.eval(field("proficiency").is_in(_PROFICIENT).cast(pl.Int32)).list.sum()
expert = pl.col("skills").list.eval((field("proficiency") == "expert").cast(pl.Int32)).list.sum()
num_skills = pl.col("flat_skills").list.len()
return [
num_skills.alias("num_skills"),
pl.col("career_history").list.len().alias("num_roles"),
durations.list.mean().fill_null(0.0).alias("avg_tenure_months"),
pl.col("career_history").list.eval(field("company").cast(pl.String)).list.n_unique().alias("num_companies"),
pl.col("flat_career_industries").list.eval(pl.element().filter(pl.element() != "")).list.n_unique().alias("num_industries"),
advanced.alias("num_advanced_skills"),
expert.alias("num_expert_skills"),
assessed.list.len().alias("num_assessments_taken"),
assessed.list.mean().alias("avg_skill_assessment"),
assessed.list.max().alias("max_skill_assessment"),
pl.when(num_skills > 0).then(assessed.list.len() / num_skills).otherwise(None).alias("assessment_coverage"),
pl.col("education").list.len().alias("num_degrees"),
pl.col("certifications").list.len().alias("num_certifications"),
(
pl.col("verified_email").cast(pl.Int32) + pl.col("verified_phone").cast(pl.Int32)
+ pl.col("linkedin_connected").cast(pl.Int32)
).alias("verification_count"),
(pl.col("notice_period_days") <= _AVAILABLE_SOON_DAYS).alias("is_available_soon"),
pl.col("career_history").list.eval(
pl.when(field("is_current")).then(field("duration_months").cast(pl.Int64)).otherwise(None)
).list.max().fill_null(0).alias("current_role_tenure_months"),
]
def _signal_exprs(as_of: date) -> list[pl.Expr]:
"""Date deltas + sentinel-aware clean signals (negative github/offer == 'no data')."""
github = pl.col("github_activity_score")
offer = pl.col("offer_acceptance_rate")
salary_min = pl.col("expected_salary_min").fill_null(0.0)
salary_max = pl.col("expected_salary_max").fill_null(0.0)
return [
((pl.lit(as_of) - to_date(pl.col("last_active_date"))).dt.total_days().fill_null(999).cast(pl.Int32)).alias("days_since_active"),
((pl.lit(as_of) - to_date(pl.col("signup_date"))).dt.total_days().fill_null(0).cast(pl.Int32)).alias("account_age_days"),
(github >= 0).alias("has_github"),
pl.when(github >= 0).then(github.clip(0, 100)).otherwise(None).alias("github_score_clean"),
(offer >= 0).alias("has_offer_history"),
pl.when(offer >= 0).then(offer.clip(0, 1)).otherwise(None).alias("offer_acceptance_rate_clean"),
salary_min.alias("salary_min"),
salary_max.alias("salary_max"),
((salary_min + salary_max) / 2.0).alias("salary_mid"),
]
def _text_section_exprs() -> list[pl.Expr]:
"""Per-section text blocks fed into the embedded body (identity / skills / experience / education)."""
field = pl.element().struct.field
industries = pl.col("flat_career_industries").list.eval(pl.element().filter(pl.element() != "")).list.unique()
identity = pl.concat_str(
[_text("headline"), _text("summary"), _text("current_title"), pl.lit(" @ "), _text("current_company")],
separator="\n",
).str.strip_chars()
core = pl.col("skills").list.eval(
pl.when(field("proficiency").is_in(_PROFICIENT)).then(field("name").cast(pl.String)).otherwise(None)
).list.drop_nulls().list.join(", ")
familiar = pl.col("skills").list.eval(
pl.when(field("proficiency").is_in(["beginner", "intermediate"])).then(field("name").cast(pl.String)).otherwise(None)
).list.drop_nulls().list.join(", ")
experience = pl.col("career_history").list.eval(
pl.concat_str([
field("title").cast(pl.String).fill_null(""), pl.lit(" @ "),
field("company").cast(pl.String).fill_null(""), pl.lit(" ("),
field("industry").cast(pl.String).fill_null(""), pl.lit(") - "),
field("description").cast(pl.String).fill_null(""),
])
).list.head(8).list.join("\n")
education = pl.col("education").list.eval(
pl.concat_str([
field("degree").cast(pl.String).fill_null(""), pl.lit(" in "),
field("field_of_study").cast(pl.String).fill_null(""), pl.lit(", "),
field("institution").cast(pl.String).fill_null(""),
])
).list.join("\n")
return [
industries.list.join(", ").alias("industries_text"),
identity.alias("identity_text"),
core.alias("core_skills_text"),
familiar.alias("familiar_skills_text"),
experience.alias("experience_text"),
education.alias("education_text"),
pl.col("flat_certifications").list.eval(pl.element().filter(pl.element() != "")).list.join(", ").alias("certifications_text"),
]
def _summary_text_expr() -> pl.Expr:
"""The deterministic, templated candidate summary — THE embedded body (no platform signals)."""
years = pl.col("years_of_experience").round(0).cast(pl.Int64).cast(pl.String)
best_tier = pl.col("flat_institution_tiers").list.eval(pl.element().filter(pl.element() != "")).list.max().fill_null("")
lead = pl.concat_str([
_text("current_title"), pl.lit(" with "), years, pl.lit(" years of experience"),
pl.when(pl.col("industries_text").str.len_chars() > 0).then(pl.lit(" in ") + pl.col("industries_text")).otherwise(pl.lit("")),
pl.lit(" ("), _text("experience_band"), pl.lit("-band)."),
])
education_clause = pl.when(pl.col("num_degrees") > 0).then(
pl.lit("Education: ") + pl.col("education_text").str.replace_all("\n", "; ")
+ pl.when(best_tier.str.len_chars() > 0).then(pl.lit(" (") + best_tier + pl.lit(" tier).")).otherwise(pl.lit("."))
).otherwise(pl.lit(""))
return pl.concat_str(
[
lead,
_clause("Currently at ", _text("current_company")),
_clause("Location: ", _text("location")), # so cities are visible to bm25 + the embedders
_clause("Core skills: ", pl.col("core_skills_text")),
_clause("Also familiar with: ", pl.col("familiar_skills_text")),
_clause("Recent roles: ", pl.col("experience_text").str.replace_all("\n", "; ")),
education_clause,
_clause("Certifications: ", pl.col("certifications_text")),
],
separator=" ", ignore_nulls=True,
).str.replace_all(r"\s+", " ").str.strip_chars()
def _clause(label: str, value: pl.Expr) -> pl.Expr:
"""Prefix ``value`` with ``label`` when non-empty, else null (dropped by concat ignore_nulls)."""
return pl.when(value.str.len_chars() > 0).then(pl.lit(label) + value).otherwise(None)
def _sha16(values: list[str]) -> list[str]:
"""sha1[:16] of each summary body — changes when the embedded text changes (re-embed trigger)."""
return [hashlib.sha1(value.encode("utf-8")).hexdigest()[:16] for value in values]
def materialize(candidates: pl.DataFrame, as_of: date) -> pl.DataFrame:
"""Add structural counts, clean signals, experience_band, and the embedded ``text`` body. Never drops."""
with step("preprocess.materialize", rows_in=candidates.height) as metrics:
# Ordered passes — text sections must exist before the summary that composes them.
out = candidates.with_columns(
_experience_band_expr(), *_count_exprs(), *_signal_exprs(as_of)
).with_columns(_text_section_exprs())
out = out.with_columns(_summary_text_expr().alias("summary_text"))
out = out.with_columns(pl.col("summary_text").alias("text"))
# Provisional hash — the next preprocess stage (quality.py) appends a strengths clause to ``text`` and
# re-hashes it there, so this value is superseded before anything reads it in the normal pipeline order.
out = out.with_columns(
pl.Series("embed_source_hash", _sha16(out.get_column("summary_text").to_list()))
)
# Single rounding pass: every derived float to 3 dp, leaving raw source signals untouched.
round_cols = [
name for name, dtype in out.schema.items()
if dtype in (pl.Float32, pl.Float64) and name not in _RAW_FLOAT_PASSTHROUGH
]
out = out.with_columns([pl.col(name).round(3) for name in round_cols])
metrics["rows_out"] = out.height
metrics["columns"] = out.width
return out