reranker / src /preprocess /validate.py
Hemprasad Badgujar
Add JD parsing, central config, checkpoints & LLM
700e896
Raw History Blame Contribute Delete
8.47 kB
"""Preprocess sub-phase 3 β€” data validation: independent integrity + graded coherence checks.
Each check is an INDEPENDENT named boolean column with a ``chk_`` prefix; there is no tiering between
them β€” they differ only in how they're consumed:
integrity checks (binary) β€” logically impossible facts; ANY hit flips ``validation_failed`` (the
candidate is treated as untrustworthy, not merely low-quality).
graded checks (weighted) β€” "smells off but possible" signals; each contributes a small penalty that
sums into a capped ``incoherence_score`` used softly downstream.
Runs as one vectorized pass on RAW input (before enrich) so flags reflect what the candidate submitted.
Flags only β€” NEVER drops a row; the hard filter that acts on these is a later phase.
"""
from __future__ import annotations
from datetime import date
import polars as pl
from common.logging import step
from .load import to_date
# Columns this stage appends to every row (the downstream contract). Order is display-only.
FLAG_COLUMNS = [
"validation_failed", "validation_rules", # integrity verdict + audit trail
"incoherence_score", "coherence_rules", # graded score + audit trail
"years_exp_claimed", "years_exp_computed", # claimed vs duration-computed YoE, side by side
"verified_skill_count", "trust_score", # quality / identity-verification counts
"date_order_ok", "salary_range_ok", "active_validated", "experience_validated", # 0/1 markers
"chk_expert_zero", "chk_career_dates_invalid", "chk_yoe_impossible", # integrity checks
"chk_date_order_invalid", "chk_salary_incoherent", # graded checks
]
# Max gap (years) between claimed and computed YoE for the claim to count as "validated" β€” 6 months.
EXPERIENCE_TOLERANCE_YEARS = 0.5
# Upper bound on summed graded penalties β€” keeps one noisy profile from dominating scoring.
COHERENCE_CAP = 0.5
# Integrity checks (binary): any True => the row is hard-failed.
_INTEGRITY_CHECKS = ["chk_expert_zero", "chk_career_dates_invalid", "chk_yoe_impossible"]
# Graded checks: (check_column, penalty weight). Heavier weight = stronger incoherence signal.
_GRADED_CHECKS = [("chk_date_order_invalid", 0.15), ("chk_salary_incoherent", 0.10)]
def _derived_exprs(as_of: date) -> list[pl.Expr]:
"""Intermediate ``d_*`` columns the checks depend on (career-history dates only).
Skill durations are intentionally NOT summed: skills run in parallel within a role, so summed
skill-months aren't commensurable with calendar career-months. Only career dates give a sound bound.
"""
field = pl.element().struct.field
earliest_start = pl.col("career_history").list.eval(to_date(field("start_date"))).list.min()
return [
# Outer bound on plausible experience; claimed YoE shouldn't exceed this by much.
((pl.lit(as_of) - earliest_start).dt.total_days() / 365.25)
.cast(pl.Float32).alias("d_career_span_years"),
# YoE implied by summing role durations (independent of the self-reported field).
(pl.col("career_history").list.eval(field("duration_months")).list.sum() / 12.0)
.cast(pl.Float32).alias("d_computed_yoe"),
]
def _check_exprs() -> list[pl.Expr]:
"""The three integrity checks plus the date-order graded signal (booleans)."""
field = pl.element().struct.field
return [
# Claims advanced/expert mastery yet logged zero months on the skill β€” contradictory.
pl.col("skills").list.eval(
field("proficiency").is_in(["advanced", "expert"]) & (field("duration_months") == 0)
).list.any().alias("chk_expert_zero"),
# A role that ends before it starts is impossible; null dates count as fine (not a hit).
pl.col("career_history").list.eval(
(to_date(field("end_date")) < to_date(field("start_date"))).fill_null(False) # noqa: E501
).list.any().alias("chk_career_dates_invalid"),
# Last-active before signup is a graded timeline glitch rather than a disqualifier.
(to_date(pl.col("last_active_date")) < to_date(pl.col("signup_date")))
.fill_null(False).alias("chk_date_order_invalid"),
# Claimed YoE exceeds the calendar span by >2yr β€” impossible (2yr slack absorbs gaps/overlap).
(pl.col("years_of_experience") > pl.col("d_career_span_years") + 2)
.fill_null(False).alias("chk_yoe_impossible"),
]
def _coherence_exprs() -> list[pl.Expr]:
"""Boolean "smells off but possible" columns; each maps to a penalty in ``_GRADED_CHECKS``."""
return [
# Inverted salary band, or a max so large it's plainly a unit/typo error (>1000 in k-units).
(
(pl.col("expected_salary_min") > pl.col("expected_salary_max"))
| (pl.col("expected_salary_max") > 1000)
).fill_null(False).alias("chk_salary_incoherent"),
]
def _audit_trail(checks: list[str], alias: str) -> pl.Expr:
"""Space-joined names of whichever of ``checks`` fired on a row β€” a human-readable audit string."""
fired = [pl.when(pl.col(check)).then(pl.lit(check)).otherwise(pl.lit("")) for check in checks]
return pl.concat_str(fired, separator=" ").str.strip_chars().alias(alias)
def _assemble_exprs() -> list[pl.Expr]:
"""Fold the raw check columns into the user-facing flags, scores, and audit-trail strings."""
# Hard-fail if ANY integrity check fired.
failed = pl.any_horizontal([pl.col(check) for check in _INTEGRITY_CHECKS]).alias("validation_failed")
# Sum each fired graded check's penalty weight, then clamp at COHERENCE_CAP.
weighted_penalty = sum(
(pl.col(check).cast(pl.Float64) * weight for check, weight in _GRADED_CHECKS), pl.lit(0.0)
)
incoherence = (
pl.min_horizontal(pl.lit(COHERENCE_CAP), weighted_penalty).cast(pl.Float32).alias("incoherence_score")
)
# Identity-verification tally (email + phone + linkedin), 0..3.
trust = (
pl.col("verified_email").cast(pl.Int8)
+ pl.col("verified_phone").cast(pl.Int8)
+ pl.col("linkedin_connected").cast(pl.Int8)
).alias("trust_score")
return [
failed,
_audit_trail(_INTEGRITY_CHECKS, "validation_rules"),
incoherence,
_audit_trail([check for check, _ in _GRADED_CHECKS], "coherence_rules"),
# Count of skills carrying an assessment-score entry β€” a quality proxy.
pl.col("verified_skills").list.len().cast(pl.Int8).alias("verified_skill_count"),
trust,
pl.col("years_of_experience").cast(pl.Float32).alias("years_exp_claimed"),
pl.col("d_computed_yoe").cast(pl.Float32).alias("years_exp_computed"),
]
def validate(candidates: pl.DataFrame, as_of: date) -> pl.DataFrame:
"""Vectorized integrity + graded coherence checks; adds the check + flag/score columns. Never drops."""
with step("preprocess.validate", rows_in=candidates.height) as metrics:
# Staged so each stage can reference the previous: derived d_* β†’ checks β†’ coherence β†’ assemble β†’ 0/1.
validated = (
candidates.with_columns(_derived_exprs(as_of))
.with_columns(_check_exprs())
.with_columns(_coherence_exprs())
.with_columns(_assemble_exprs())
.with_columns(
date_order_ok=(~pl.col("chk_date_order_invalid")).cast(pl.Int8),
salary_range_ok=(~pl.col("chk_salary_incoherent")).cast(pl.Int8),
)
.with_columns(
# 1 iff signup precedes/equals last-active; missing dates fall to 0 (stricter than the marker).
active_validated=(to_date(pl.col("signup_date")) <= to_date(pl.col("last_active_date")))
.fill_null(False).cast(pl.Int8),
# 1 iff claimed YoE agrees with the duration-summed computed YoE to within 6 months.
experience_validated=(
(pl.col("years_of_experience") - pl.col("d_computed_yoe")).abs()
<= EXPERIENCE_TOLERANCE_YEARS
).fill_null(False).cast(pl.Int8),
)
)
failed_count = int(validated.select(pl.col("validation_failed").sum()).item())
metrics["rows_out"] = validated.height
metrics["validation_failed"] = failed_count
metrics["invalid_rate"] = round(failed_count / validated.height, 5) if validated.height else 0.0
return validated