Files
VacatureRadar/apps/jobs/services/pipeline.py
T
Jens cbd7220d8d
deploy / deploy (push) Canceled after 0s
perf: optimize import persistence for 0.3.14
2026-07-29 18:56:11 +02:00

489 lines
17 KiB
Python

from __future__ import annotations
from dataclasses import dataclass
from dataclasses import field as dataclass_field
from decimal import Decimal
from urllib.parse import urlsplit
from django.db import IntegrityError, transaction
from django.utils import timezone
from apps.jobs.models import (
Employer,
FieldProvenance,
JobPosting,
JobSourceAlias,
JobVersion,
ScoreRun,
)
from apps.profiles.models import SearchProfile
from apps.sources.adapters.base import FieldEvidence
from apps.sources.adapters.registry import registry
from apps.sources.models import RawDocument, Source
from apps.sources.services.policy import is_denied_domain
from .dedupe import DedupeDecision, find_existing_job
from .features import extract_deterministic_features
from .normalization import CanonicalJobDraft, normalize_extracted_job, normalize_token
from .scoring import score_and_save
RECRUITER_TERMS = {
"recruitment",
"recruiter",
"staffing",
"interim",
"consultancy",
"consulting",
"talent",
}
@dataclass
class PersistenceContext:
"""Run-scoped reference cache; never shared between workers or imports."""
employers: dict[tuple[str, str], Employer | None] = dataclass_field(default_factory=dict)
active_profiles: list[SearchProfile] | None = None
latest_scores: dict[tuple[object, int], ScoreRun] = dataclass_field(default_factory=dict)
def _confidence(value: float) -> Decimal:
return Decimal(str(max(0.0, min(1.0, value))))
def resolve_employer(
draft: CanonicalJobDraft, *, context: PersistenceContext | None = None
) -> Employer | None:
name = draft.employer_name.strip()
domain = draft.employer_domain.strip()
if not name and not domain:
return None
display_name = name or domain
normalized = normalize_token(display_name)
cache_key = (normalized, domain)
if context is not None and cache_key in context.employers:
return context.employers[cache_key]
recruiter = any(term in normalized for term in RECRUITER_TERMS)
employer, _ = Employer.objects.get_or_create(
normalized_name=normalized,
domain=domain,
defaults={
"name": display_name,
"is_direct_employer": not recruiter,
"is_recruiter": recruiter,
"confidence": _confidence(0.75 if name else 0.45),
},
)
changed: list[str] = []
if name and employer.name != name and len(name) > len(employer.name):
employer.name = name
changed.append("name")
if recruiter and not employer.is_recruiter:
employer.is_recruiter = True
employer.is_direct_employer = False
changed.extend(["is_recruiter", "is_direct_employer"])
if changed:
employer.save(update_fields=[*changed, "updated_at"])
if context is not None:
context.employers[cache_key] = employer
return employer
def job_snapshot(job: JobPosting) -> dict[str, object]:
return {
"id": str(job.id),
"title": job.original_title,
"normalized_title": job.normalized_title,
"employer": job.employer_name,
"canonical_url": job.canonical_url,
"location": job.raw_location,
"region": job.region,
"municipality": job.municipality,
"workplace_type": job.workplace_type,
"employment_types": job.employment_types,
"description_text": job.description_text,
"skills_required": job.skills_required,
"skills_preferred": job.skills_preferred,
"date_posted": job.date_posted.isoformat() if job.date_posted else None,
"valid_through": job.valid_through.isoformat() if job.valid_through else None,
"status": job.status,
"content_hash": job.content_hash,
}
def _source_is_direct(source: Source | None, draft: CanonicalJobDraft) -> bool:
host = (urlsplit(draft.canonical_url).hostname or "").lower()
if is_denied_domain(host):
return False
if not source:
return False
metadata = source.metadata if isinstance(source.metadata, dict) else {}
return source.source_type == Source.Type.EMPLOYER or metadata.get("direct_employer") is True
def _alias_payload(
raw_payload: object,
decision: DedupeDecision,
*,
fallback_canonical_url: str | None,
) -> dict[str, object]:
if not isinstance(raw_payload, dict):
raw_payload = {}
payload: dict[str, object] = dict(raw_payload)
if decision.reason == "review_direct_conflict" or decision.resolved_direct:
payload["employer_resolution"] = {
"reason": decision.reason,
"confidence": float(decision.similarity),
"canonical_url": decision.canonical_url or fallback_canonical_url,
"resolved_direct": decision.resolved_direct,
"evidence": decision.evidence,
}
return payload
def _apply_draft(
job: JobPosting,
draft: CanonicalJobDraft,
employer: Employer | None,
*,
direct: bool,
resolved_canonical_url: str | None = None,
) -> list[str]:
canonical_url = resolved_canonical_url if direct else draft.canonical_url
fields = {
"employer": employer,
"original_title": draft.title,
"normalized_title": draft.normalized_title,
"job_family": draft.job_family,
"language": draft.language,
"description_html_sanitized": draft.description_html,
"description_text": draft.description_text,
"raw_location": draft.location_text,
"country": draft.country,
"region": draft.region,
"municipality": draft.municipality,
"postal_code": draft.postal_code,
"workplace_type": draft.workplace_type,
"employment_types": draft.employment_types,
"compensation": draft.compensation,
"skills_required": draft.skills_required,
"skills_preferred": draft.skills_preferred,
"date_posted": draft.date_posted,
"valid_through": draft.valid_through,
"content_hash": draft.content_hash,
"analysis_features": extract_deterministic_features(draft.title, draft.description_text),
"status": JobPosting.Status.ACTIVE,
"last_seen": timezone.now(),
"direct_employer": direct or (employer.is_direct_employer if employer else False),
"recruiter": employer.is_recruiter if employer else not direct,
}
changed: list[str] = []
for field, value in fields.items():
if (value not in (None, "", [], {}) or field in {"status", "last_seen"}) and (
getattr(job, field) != value
):
setattr(job, field, value)
changed.append(field)
if direct and canonical_url and job.canonical_url != canonical_url:
job.canonical_url = canonical_url
changed.append("canonical_url")
substantive_changes = [field for field in changed if field != "last_seen"]
if substantive_changes and "last_changed" not in changed:
job.last_changed = timezone.now()
changed.append("last_changed")
return changed
def _sync_evidence(
*,
job: JobPosting,
alias: JobSourceAlias,
evidence_items: list[FieldEvidence],
parser_version: str,
) -> None:
existing = {
(item.field_name, item.extraction_method): item
for item in FieldProvenance.objects.filter(job=job, source_alias=alias)
}
creates: list[FieldProvenance] = []
updates: list[FieldProvenance] = []
for evidence in evidence_items:
key = (evidence.field_name, evidence.method)
values = {
"confidence": _confidence(evidence.confidence),
"evidence_excerpt": evidence.evidence[:1000],
"parser_version": parser_version,
}
current = existing.get(key)
if current is None:
creates.append(
FieldProvenance(
job=job,
source_alias=alias,
field_name=evidence.field_name,
extraction_method=evidence.method,
**values,
)
)
continue
changed = False
for field_name, value in values.items():
if getattr(current, field_name) != value:
setattr(current, field_name, value)
changed = True
if changed:
current.updated_at = timezone.now()
updates.append(current)
if creates:
FieldProvenance.objects.bulk_create(creates, batch_size=250)
if updates:
FieldProvenance.objects.bulk_update(
updates,
["confidence", "evidence_excerpt", "parser_version", "updated_at"],
batch_size=250,
)
def _copy_score(score: ScoreRun) -> ScoreRun:
return ScoreRun.objects.create(
job=score.job,
profile=score.profile,
profile_version=score.profile_version,
score=score.score,
confidence=score.confidence,
recommendation=score.recommendation,
components=score.components,
positives=score.positives,
concerns=score.concerns,
hard_exclusions=score.hard_exclusions,
evidence=score.evidence,
model_version=score.model_version,
prompt_version=score.prompt_version,
)
@transaction.atomic
def persist_draft(
draft: CanonicalJobDraft,
*,
document: RawDocument,
parser_key: str,
parser_version: str,
extraction_confidence: float,
context: PersistenceContext | None = None,
) -> tuple[JobPosting, DedupeDecision, bool]:
source = document.source
employer = resolve_employer(draft, context=context)
decision = find_existing_job(draft, source=source)
direct = _source_is_direct(source, draft)
if decision.resolved_direct:
direct = True
created = False
substantive_change = False
if decision.job is None:
try:
job = JobPosting.objects.create(
employer=employer,
original_title=draft.title,
normalized_title=draft.normalized_title,
job_family=draft.job_family,
language=draft.language,
canonical_url=draft.canonical_url,
canonical_key=draft.canonical_key,
content_hash=draft.content_hash,
description_html_sanitized=draft.description_html,
description_text=draft.description_text,
raw_location=draft.location_text,
country=draft.country,
region=draft.region,
municipality=draft.municipality,
postal_code=draft.postal_code,
workplace_type=draft.workplace_type,
employment_types=draft.employment_types,
compensation=draft.compensation,
skills_required=draft.skills_required,
skills_preferred=draft.skills_preferred,
date_posted=draft.date_posted,
valid_through=draft.valid_through,
direct_employer=direct or (employer.is_direct_employer if employer else False),
recruiter=employer.is_recruiter if employer else not direct,
extraction_confidence=_confidence(extraction_confidence),
analysis_features=extract_deterministic_features(
draft.title, draft.description_text
),
status=JobPosting.Status.ACTIVE,
)
created = True
except IntegrityError:
job = JobPosting.objects.get(canonical_key=draft.canonical_key)
decision = DedupeDecision(job, "canonical_key_race", 1.0)
else:
job = decision.job
if job.content_hash != draft.content_hash:
JobVersion.objects.get_or_create(
job=job,
content_hash=job.content_hash,
defaults={"snapshot": job_snapshot(job), "changed_fields": []},
)
changed = _apply_draft(
job,
draft,
employer,
direct=direct,
resolved_canonical_url=decision.canonical_url,
)
if extraction_confidence > float(job.extraction_confidence):
job.extraction_confidence = _confidence(extraction_confidence)
changed.append("extraction_confidence")
substantive_change = any(field not in {"last_seen", "last_changed"} for field in changed)
if changed:
job.save(update_fields=list(dict.fromkeys([*changed, "updated_at"])))
alias = (
JobSourceAlias.objects.filter(
job=job,
source=source,
canonical_url=draft.canonical_url,
external_id=draft.external_id,
)
.order_by("-last_seen")
.first()
)
if alias is None:
alias = JobSourceAlias.objects.create(
job=job,
source=source,
raw_document=document,
url=draft.source_url,
canonical_url=draft.canonical_url,
external_id=draft.external_id,
source_title=draft.title,
source_employer=draft.employer_name,
extraction_method=parser_key,
extraction_confidence=_confidence(extraction_confidence),
is_canonical=direct,
payload=_alias_payload(
draft.raw,
decision,
fallback_canonical_url=draft.canonical_url,
),
)
else:
alias.last_seen = timezone.now()
alias.raw_document = document
alias.payload = _alias_payload(
alias.payload, decision, fallback_canonical_url=draft.canonical_url
)
if direct:
alias.is_canonical = True
alias.save(
update_fields=["last_seen", "raw_document", "payload", "is_canonical", "updated_at"]
)
_sync_evidence(
job=job,
alias=alias,
evidence_items=draft.evidence,
parser_version=parser_version,
)
if created or substantive_change:
JobVersion.objects.get_or_create(
job=job,
content_hash=job.content_hash,
defaults={"snapshot": job_snapshot(job), "changed_fields": []},
)
if context is not None:
if context.active_profiles is None:
context.active_profiles = list(SearchProfile.objects.filter(is_active=True))
profiles = context.active_profiles
else:
profiles = SearchProfile.objects.filter(is_active=True)
for profile in profiles:
score = score_and_save(job, profile)
if context is not None:
context.latest_scores[(job.pk, profile.pk)] = score
else:
if context is not None:
if context.active_profiles is None:
context.active_profiles = list(SearchProfile.objects.filter(is_active=True))
profiles = context.active_profiles
else:
profiles = SearchProfile.objects.filter(is_active=True)
for profile in profiles:
cache_key = (job.pk, profile.pk)
previous = context.latest_scores.get(cache_key) if context is not None else None
if previous is None:
previous = (
ScoreRun.objects.filter(job=job, profile=profile)
.order_by("-created_at")
.first()
)
if previous is None or previous.profile_version != profile.version:
score = score_and_save(job, profile)
else:
score = _copy_score(previous)
if context is not None:
context.latest_scores[cache_key] = score
return job, decision, created
@transaction.atomic
def process_raw_document(
document: RawDocument, *, context: PersistenceContext | None = None
) -> dict[str, int | str | list[str]]:
result = registry.extract(document)
document.parser_key = result.parser_key
document.parser_version = result.parser_version
document.extraction_confidence = _confidence(result.confidence)
document.save(
update_fields=["parser_key", "parser_version", "extraction_confidence", "updated_at"]
)
created = updated = duplicates = 0
for extracted in result.jobs:
source_metadata = (
document.source.metadata
if document.source and isinstance(document.source.metadata, dict)
else {}
)
configured_employer = str(source_metadata.get("employer_name") or "").strip()
if configured_employer:
extracted.employer_name = configured_employer
extracted.evidence.append(
FieldEvidence(
"employer_name",
"reviewed-source-config",
0.99,
configured_employer,
)
)
draft = normalize_extracted_job(extracted)
_, decision, was_created = persist_draft(
draft,
document=document,
parser_key=result.parser_key,
parser_version=result.parser_version,
extraction_confidence=result.confidence,
context=context,
)
if was_created:
created += 1
elif (
decision.reason.startswith("exact")
or decision.reason.startswith("fuzzy")
or decision.reason == "resolved_direct_match"
):
duplicates += 1
updated += 1
else:
updated += 1
return {
"extracted": len(result.jobs),
"created": created,
"updated": updated,
"duplicates": duplicates,
"parser": result.parser_key,
"warnings": result.warnings,
}