File size: 8,466 Bytes
700e896 | 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 | """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
|