Clinical Notes Pipeline
Overview
The ingestion pipeline loads data in two steps between a hospital's FHIR server and the Cavell extraction API:
- Seed — organizations, practitioners, and patients into FHIR (enforces correct reference ordering internally)
- Extract — clinical documents, processed per patient in chronological order. Use
extract_all()for a whole dataset (globally date-sorted, batched) orextract()for a single pass; see Extract Options.
The pipeline resolves identifiers to IDs, so you work with business identifiers (patient IDs, facility codes, staff IDs) rather than FHIR server-assigned IDs.
Before each extraction, the SDK passes the patient, organization, and practitioner FHIR IDs as explicit API params, alongside any existing clinical resources as context. The API uses these real FHIR references on all extracted resources — no reference rewriting needed.
Practitioners are matched by the SDK after extraction — the Cavell API extracts practitioner names from text, and the SDK links them to seeded practitioners in FHIR.
This page covers clinical notes, which are extracted by a model. Already- structured lab results take a separate, deterministic path with no LLM and no token cost — see the lab results pipeline.
Data Flow
Seed: organizations + practitioners + patients → FHIR server
Extract: For each document:
1. Fetch existing clinical resources from FHIR
2. Send text + clinical context + reference IDs as params → Cavell API
3. Match extracted practitioners to seeded ones
4. Persist clinical resources with correct references → FHIR server
Start from a CSV
Most work starts from an export rather than hand-built objects. All three
data types provide a from_rows() classmethod that builds objects from CSV data (or any list of dicts). The columns dict maps SDK field names (left) to your CSV column headers (right). Keyword arguments are applied as literal defaults to every row.
import csv
from cavell_client import Patient, Practitioner, Document
with open("notes.csv", newline="", encoding="utf-8") as f:
rows = list(csv.DictReader(f))
practitioners = Practitioner.from_rows(
rows,
columns={"identifier": "practitioner_id", "name": "practitioner_name"},
organization_identifier="DEMO-HOSPITAL",
)
patients = Patient.from_rows(
rows,
columns={
"identifier": "patient_id",
"name": "patient_name",
"birth_date": "birth_date",
"gender": "gender",
"general_practitioners": "practitioner_id",
},
managing_organization="DEMO-HOSPITAL",
)
documents = Document.from_rows(
rows,
columns={
"text": "note_text",
"patient_identifier": "patient_id",
"date": "note_date",
"document_id": "note_id",
"encounter_id": "encounter_id",
"practitioner_identifier": "practitioner_id",
"meta": "department",
},
organization_identifier="DEMO-HOSPITAL",
)
Patient.from_rows() and Practitioner.from_rows() deduplicate by identifier (first occurrence wins) and skip rows with empty identifiers. Document.from_rows() creates one document per row, validates that document_id values are unique, and warns on short text (< 20 characters). Its columns mapping must include text, patient_identifier, date and document_id; a row with a blank value for any of them raises rather than being silently coerced to None.
Document.from_rows() preserves CSV row order and does not sort by date. To
ingest the whole file, hand the list to
extract_all(), which sorts globally by date
before batching:
pipeline.seed(organizations=[...], patients=patients, practitioners=practitioners)
outcomes = pipeline.extract_all(documents, batch_size=500)
Practitioner.from_rows() supports a virtual "name" column that auto-splits "Given Family" into given_name and family_name. Use either "name" or the individual "given_name"/"family_name" keys, not both.
All three validate upfront — unknown field names, missing CSV columns, and missing required keys raise ValueError before any objects are built.
Every column's meaning, and what it produces downstream, is in the CSV field reference.
Full Walkthrough
The same pipeline with every object built by hand — every option in one place.
from cavell_client import (
CavellClient,
IngestionPipeline,
Organization,
Practitioner,
Patient,
Document,
)
with CavellClient(
api_url="https://prd.prism.cavell.app/api",
api_key="your-llm-gateway-key",
fhir_base_url="http://localhost:8090",
fhir_client_id="...", # optional; provide together with fhir_client_secret
fhir_client_secret="...", # optional; omit both for unauthenticated FHIR servers
) as client:
# Connectivity is already verified: construction raises if the key is
# rejected or either endpoint is unreachable.
# Create pipeline with optional tier selection and concurrency
pipeline = IngestionPipeline(
client,
tier="low",
max_concurrency=10,
default_organization="CGH-001",
)
# Seed organizations, practitioners, and patients
pipeline.seed(
organizations=[
Organization(identifier="CGH-001", name="City General Hospital"),
Organization(identifier="SMH-002", name="St. Mary's Hospital"),
],
patients=[
Patient(
identifier="MRN-12345",
name="John Doe",
birth_date="1985-03-15",
gender="male",
managing_organization="CGH-001",
general_practitioners=["DOC-001"],
),
Patient(
identifier="MRN-67890",
managing_organization="SMH-002",
),
],
practitioners=[
Practitioner(
identifier="DOC-001",
family_name="Smith",
given_name="Jane",
organization_identifier="CGH-001",
),
Practitioner(
identifier="DOC-002",
family_name="Jones",
given_name="Bob",
organization_identifier="SMH-002",
),
],
)
# Extract documents (skip_processed=True by default, so reruns are safe
# when documents have stable document_id values)
outcomes = pipeline.extract(
[
Document(
text="Patient presents with type 2 diabetes...",
patient_identifier="MRN-12345",
date="2024-01-15",
practitioner_identifier="DOC-001",
document_id="note-001",
),
Document(
text="Follow-up: diabetes well controlled on metformin...",
patient_identifier="MRN-12345",
date="2024-04-15",
document_id="note-002",
),
Document(
text="Patient reports chest pain...",
patient_identifier="MRN-67890",
date="2024-02-01",
document_id="note-003",
organization_identifier="SMH-002", # overrides default
),
]
)
for outcome in outcomes:
if outcome.success:
print(
f"[{outcome.patient_identifier}] "
f"Extracted {outcome.extract_result.count} resources"
)
else:
print(f"[{outcome.patient_identifier}] Error: {outcome.error}")
Helper Dataclasses
Organization
| Field | Type | Required | Description |
|---|---|---|---|
identifier |
str |
Yes | Facility code (e.g., "CGH-001") |
name |
str |
Yes | Display name for the extraction API |
Practitioner
| Field | Type | Required | Description |
|---|---|---|---|
identifier |
str |
Yes | Staff ID (e.g., "DOC-001") |
family_name |
str |
Yes | Family name |
given_name |
str |
Yes | Given name |
organization_identifier |
str |
Yes | Organization identifier (must match a seeded org) |
specialty |
str |
No | Clinical specialty (written to PractitionerRole.specialty) |
Patient
| Field | Type | Required | Description |
|---|---|---|---|
identifier |
str |
Yes | Patient identifier (e.g., MRN) |
name |
str |
No | Patient name |
birth_date |
str |
No | ISO date (e.g., "1990-01-15") |
gender |
str |
No | FHIR gender code |
managing_organization |
str |
No | Organization identifier |
general_practitioners |
str or list[str] |
No | Practitioner identifier(s); a single string is normalized to a list |
Document
| Field | Type | Required | Description |
|---|---|---|---|
text |
str |
Yes | Clinical text |
patient_identifier |
str |
Yes | Patient identifier (e.g., MRN) |
date |
str or date |
Yes | ISO date ("2024-01-15") or datetime.date, normalized to YYYY-MM-DD on construction. Drives chronological ordering, and is sent to the API as the document_date payload field |
document_id |
str |
Yes | Your identifier for this document, stamped on the DocumentReference. Keyword-only. Everything that makes re-running safe keys on it — the resume filter, the chronology watermark, and failure reporting |
organization_identifier |
str |
No | Org identifier — falls back to default_organization if omitted |
meta |
str |
No | Extra context for the extraction API (e.g. department, ward). Do not include the document date or practitioner — the date is sent as its own document_date payload field and the practitioner is injected automatically (see Meta Assembly). |
practitioner_identifier |
str |
No | If provided, the SDK resolves this to a FHIR ID (passed as a param) and injects the practitioner's name into meta, improving matching precision |
encounter_id |
str |
No | Your identifier for the visit/admission the document belongs to. When set, the SDK looks the Encounter up in FHIR (urn:cavell:encounter) and sends it to the API for an in-place update; the API creates it when there is none yet, and links every resource extracted from the document to it. When unset, no Encounter is created for the document |
Field validation and normalization happen in one place: building a Document
is what checks it. extract() and extract_all() add only a type check —
passing raw CSV rows straight through raises TypeError naming the offending
positions, rather than failing inside a worker thread mid-run. It runs before
the API pre-flight, so a malformed list costs not even one request.
IngestionOutcome
Returned by pipeline.extract() for each document.
| Field | Type | Description |
|---|---|---|
success |
bool |
Whether extraction and persistence succeeded |
patient_identifier |
str |
Patient identifier from the document |
document_index |
int |
Position in the list extract() actually processed — i.e. after skip_processed filtering and the limit truncation, not in the list you passed. Under extract_all() it is relative to the document's own batch |
document_id |
str |
The document_id from the source Document — the stable way to correlate an outcome back to a source record |
extract_result |
ExtractResult or None |
Result if successful |
error |
str or None |
Error message if failed |
transient |
bool |
The failure was transient (timeout, connection drop, server 5xx/429) — the deferred retry pass re-ran or will re-run it (see Document extraction failures) |
out_of_order |
bool |
The document predated its patient's newest already-persisted document, so it was extracted against split context. This records how the document was processed, not whether it worked — it pairs with success=True on a normal run |
Note on naming: Document.meta is the SDK's free-text supplementary context
for the extraction model. It is unrelated to FHIR Resource.meta, the
server-side element that carries version, provenance, and validation tags.
Practitioner Matching
After extraction, the SDK matches practitioners found in the text to those seeded via seed():
- 1 match — the reference is linked to the seeded practitioner.
- 0 matches — the reference is removed and a warning is logged.
- >1 matches — the match is treated as ambiguous, the reference is removed, and a warning is logged.
Matching uses identifier first (exact match), then falls back to family name + given name scoped to the document's organization.
Practitioner resources are always removed from the persisted bundle — only clinical resources (Conditions, MedicationRequests, etc.) are written to FHIR with correct practitioner references.
Ordering and Dates
Every document requires a date, and documents reach the API in ascending date order per patient. This ordering determines correct behavior -- see How Updates Work.
Sorting happens at two levels:
extract_all()sorts the entire dataset by ascending date before batching, so batch boundaries are chronologically clean.extract()sorts within each patient in the pass it was given. This is the safety net for callingextract()directly, and a no-op on inputextract_all()has already ordered.
Both sorts are stable, so same-day documents keep the order you passed them in.
Neither sorts your input list in place, and neither sorts the dataset when
you call extract() yourself — if you batch by hand with extract(..., limit=N),
ordering across those calls is yours to get right. Use extract_all() and it
is handled.
Out-of-order documents get split context
Sorting only orders the documents within one run. If a document is older than data already persisted for that patient (from an earlier run), you are going backwards in time. The pipeline extracts it anyway, but it changes what the document is shown.
Each document is compared against its own patient's watermark — the date of the newest already-processed document for that patient. A document older than that watermark takes the split-context path:
contextholds the patient's record as it stood on that document's date, so the API can treat the note as the latest one and reason the way its author did.future_contextholds everything already on record that the document could not have known, sent separately rather than mixed in.out_of_order: truetells the API which of the two it is looking at.
Resources are assigned by provenance, not clinical date: each is dated by
the newest already-processed document that created or updated it, read from that
document's DocumentReference.context.related. What matters is when a fact
entered the record, not when it happened — a Condition recorded last year with a
1998 onset is still knowledge this document's author could not have had. A
resource counts as past only when every document that touched it is at or
before the document's date; anything a newer document has since modified carries
that newer knowledge in its body and goes to future_context.
outcomes = pipeline.extract_all(documents, batch_size=500)
n_split = sum(1 for o in outcomes if o.out_of_order)
print(f"{n_split} document(s) extracted against split context")
IngestionOutcome.out_of_order records how a document was processed, not
whether it worked — it pairs with success=True on a normal run.
Caveats:
- Dates are truncated to the day, so two same-day documents have no defined order and neither is out of order. Process same-day documents in one run (the input order is preserved for equal dates).
- The check fails open per patient: if the watermark query errors, that patient's documents are all treated as in-order with a warning rather than failing the run.
- Documents already processed are filtered out by
skip_processedbefore the check. - The split costs one extra FHIR query per document, for provenance. Sorting your
dataset first avoids paying it at all, which is what
extract_all()does. - Observations are bounded at the document's date, and none dated after it
are sent. The API matches an observation on an exact
(date, code, value), and a document only reports results at or before its own date, so a later one could never match anything it proposes. A result measured before the document but only recorded by a later note still reachesfuture_context— provenance puts it there, not its clinical date. - ResearchStudy is exempt and always travels as ordinary context. Studies
carry no date to split on (the API writes only identifier, status, title and
arms), and a study the extractor cannot see is one it re-creates under a fresh
id — wrong for every patient on the server. Studies are capped at
MAX_CONTEXT_STUDIES(200) instead.
Requires matching API support
future_context and out_of_order must be understood by the extraction
API you are pointed at. An older deployment ignores unknown fields
silently, which would leave a backdated document extracting against
past-only context with nothing to reconcile against.
Update guard
An earlier design guarded the persist step instead, dropping updates to
resources sourced from newer documents. That is now the API's decision —
it is told which resources postdate the document and reconciles them
itself, and a blanket guard would veto exactly that. The code is retained
but disabled (_apply_update_guard, and the commented-out block in
_process_single_document) as a fallback.
Meta Assembly
Before each extraction call, the pipeline builds the meta string sent to the API by combining up to two parts in order:
- Your
metavalue (if set) — e.g.Department: Cardiology - Attending practitioner (if
practitioner_identifieris set) — e.g.Attending: Jane Smith (DOC-001)
For example, a document with meta="Department: Cardiology" and practitioner_identifier="DOC-001" produces:
When there is nothing to say — no meta, no practitioner — the field is omitted from the payload entirely rather than sent empty.
The document date is not part of meta. It travels as its own document_date payload field, normalized to YYYY-MM-DD:
{
"text": "...",
"document_date": "2024-01-15",
"meta": "Department: Cardiology\nAttending: Jane Smith (DOC-001)",
"document_identifier": "note-001"
}
Because the pipeline sends the date separately and injects the practitioner, do not duplicate either in your meta field — that would confuse the extraction model.
Error Handling
Seed failures (fail-fast)
If seeding fails, the pipeline raises RuntimeError immediately. Resources already written to the FHIR server stay there, but re-running is safe — seeding uses upsert, so it picks up where it left off.
try:
pipeline.seed(organizations=[...], patients=[...])
except RuntimeError as e:
print(f"Seeding failed: {e}")
Cross-validation failures
If you reference an unknown identifier, seed() raises ValueError before making any FHIR calls:
# This raises ValueError because "UNKNOWN-ORG" was never provided
pipeline.seed(
organizations=[Organization(identifier="CGH-001", name="City General")],
patients=[Patient(identifier="MRN-1", managing_organization="UNKNOWN-ORG")],
)
Document extraction failures
When extraction fails for a document, the pipeline:
- Returns a failed
IngestionOutcomefor that document. - Continues with the patient's remaining documents — a failed document persists nothing, so the record they extract against stays consistent.
- Continues processing other patients.
Re-ingest the failed document later (a re-run with skip_processed=True
picks it up automatically): by then it is older than the patient's newest
persisted document, so it takes the split-context path
and is reconciled against everything recorded after it.
for outcome in pipeline.extract(documents):
if not outcome.success:
print(f"Failed: {outcome.error}")
# Other patients' documents still process
Transient failures are retried. Timeouts, connection drops, and server
5xx/429 responses are retried in place (3 attempts with backoff). If a
document still fails transiently, its outcome is marked transient=True and,
after all patients finish, a deferred pass re-runs each transiently
failed document once, in date order. Context is re-fetched fresh, so this is
safe — a failed document persists nothing — and a document whose successors
persisted in the meantime takes the split-context path. Deterministic
failures (rejected bundles, 4xx content errors) are not retried: re-running
reproduces them.
Run aborts
Two failures are run-global — no document can succeed until they are fixed —
so instead of failing every document one by one, extract() raises:
CavellAuthError(401): the key is missing, expired, or rejected. Checked once cheaply before any processing; a mid-run 401 stops all patients after the first failure.CavellGatewayUnavailableError(503): the LLM Gateway is unreachable from the server. In-place retries still run (a blip should not abort), but if a document exhausts them the run stops.
Aborting is safe: failed documents persist nothing, and re-running
extract() with skip_processed=True (the default) resumes where the run
stopped.
Persistence failures
The pipeline treats these the same as extraction failures: the document gets a failed outcome and the patient's remaining documents continue.
How Updates Work
Before each document, the SDK fetches existing clinical resources from the FHIR server and passes them as context to the extraction API. The API uses this context to POST a new resource or PUT an update to an existing one.
The context covers every resource type the extraction API can read, so that each is updated in place rather than duplicated:
| Resource | Scope sent as context |
|---|---|
Condition |
all |
AllergyIntolerance |
all |
Procedure |
all |
MedicationRequest |
all |
MedicationAdministration |
all — so a dose series described across notes closes its period instead of starting a second one |
NutritionOrder |
all, including revoked/completed — a diet order can only be stopped by updating the order that started it, which requires seeing it |
FamilyMemberHistory |
all — family history is restated at nearly every visit; conditions are merged onto the relative already on record |
Observation |
the last 2 years relative to the document's own date, capped at the 50 most recent within that window |
CarePlan |
active only — a plan that ends is marked completed/revoked and drops out, so it stops growing the context |
ResearchSubject |
all |
ResearchStudy |
all statuses except entered-in-error/withdrawn. Not patient-scoped: a study applies to every patient |
Two notes on the scoping:
- The Observation window is anchored on the document, not on today. Backfilling an archive of 2015 notes with a today-anchored window would put every stored observation outside it and hand the extractor nothing.
- ResearchStudy is deliberately not filtered to
active. A study the patient joined two years ago is still the study a new note names, but its status moves on tocompletedorclosed-to-accrual. Filtering to active dropped it from the context exactly then, and the extractor — which matches studies by title and embedding similarity — re-created it under a fresh id.
Identity resources are never sent as context: Patient, Organization and Practitioner travel as explicit reference IDs (patient_id, organization_id, practitioner_id) instead. DocumentReference is note-scoped by design — one per document — so it is not context either.
Encounter is the exception, and it is fetched on its own terms rather than as part of the patient-wide context. A document with an encounter_id triggers one targeted search — Encounter?subject=<patient>&identifier=urn:cavell:encounter|<encounter_id> — and the match, if any, is appended to context so the API updates that Encounter in place (its id never changes, so nothing already linked to it loses its reference). It always travels in context, even for a backdated document: the API's out-of-order gate drops proposed updates whose id appears in future_context, and the Encounter is meant to be updated. With no match the API creates the Encounter; with no encounter_id it creates none.
Requires matching API support
An extraction API that predates encounter_identifier ignores the field
silently and produces no Encounter at all — its Encounters were decided by
the planning model, per document, and are not what this SDK now expects.
Chronological ordering matters:
- Doc 1 mentions "type 2 diabetes" → API creates a new Condition resource
- SDK persists the Condition to FHIR
- Before Doc 2, SDK fetches context again -- the Condition now exists
- Doc 2 mentions "diabetes well controlled" → API sees the existing Condition and updates it rather than creating a duplicate
Validation tags and updates
Every resource the extraction API emits carries an unvalidated meta tag
(see Mark a resource validated). By
the server contract, any update to a resource re-adds the tag — validation
applies to a specific version.
One client-side behavior softens this in practice:
- Duplicate suppression: observations that already exist (same date, code, and value) are dropped client-side before persisting, so re-extracting the same data does not touch validated resources.
An in-order update with genuinely new information does replace the resource
and re-marks it unvalidated — by design, since a clinician has not seen the
new version.
Backdated documents are not exempt from this. An out-of-order document is
extracted like any other (see Out-of-order documents get split
context), and the bundle it produces
can update a resource that a newer document created — which re-adds
unvalidated to it. The client does not veto those updates: it tells the API
which resources postdate the document (future_context) and the API's
reconciliation decides what to merge, drop or create. A client-side update guard
that dropped every PUT against a newer-sourced resource is retained but disabled
in _process_single_document, because under the split rule every such resource
postdates the document by definition, so the guard would veto exactly the
reconciliation the API was asked to make. If you need the stricter behaviour,
that is the block to re-enable.
Connection Check
You do not need to check the connection — CavellClient already did. Constructing it validates both halves of your configuration and raises on the first failure, so nothing is left to verify before seed() or extract():
| Check | Request | Raises on failure |
|---|---|---|
| Cavell API key | GET /key/info — pre-flights your LLM Gateway key against the gateway, no tokens spent |
CavellAuthError (rejected key), CavellGatewayUnavailableError (gateway down), CavellAPIError (URL doesn't serve the route) |
| FHIR server | GET /metadata |
FHIRAuthError (OAuth2 handshake failed), FHIRConnectionError (unreachable or wrong URL) |
The exception type tells you which half is misconfigured:
CavellClient(api_url=..., api_key="sk-typo", fhir_base_url=...) # pragma: allowlist secret
# CavellAuthError: Cavell API error (401): LLM Gateway rejected the provided key.
CavellClient(api_url=..., api_key=<valid>, fhir_base_url="http://localhost:9999")
# FHIRConnectionError: FHIR server unreachable: http://localhost:9999 (...)
Neither check can be skipped: the Prism API is always remote, so a configuration that can't be validated is an error rather than a supported mode. extract() re-runs the key pre-flight once per run, so a key that expires between runs never reaches the pipeline.
client.check_connection() remains available for re-checking a long-lived client (say, before a new batch hours later). Unlike the constructor it reports rather than raises, and always checks both services even if one fails:
Each key contains ok (bool) and, on failure, error (str).
Pipeline Options
pipeline = IngestionPipeline(
client,
tier="low",
max_concurrency=10,
default_organization="CGH-001",
)
| Option | Default | Description |
|---|---|---|
tier |
None |
Model tier to use for extraction (low/medium/high) |
max_concurrency |
5 |
How many patients are processed in parallel |
default_organization |
None |
Org identifier used when Document.organization_identifier is omitted |
Extract Options
Two entry points, and the difference matters for large datasets:
| Method | Scope | Ordering |
|---|---|---|
extract_all(documents, batch_size=N) |
Every document, in batches of N |
Sorts the whole dataset by ascending date first |
extract(documents, limit=N) |
One pass, capped at N documents |
Sorts by date within each patient in that pass |
extract_all() — the whole dataset
Use this for anything you would describe as "ingest this CSV". It sorts every
document by ascending date, splits the sorted list into batches, and calls
extract() once per batch.
outcomes = pipeline.extract_all(
documents,
batch_size=500, # documents per extract() call; None = one call
skip_processed=True, # applied per batch
on_batch=lambda batch: print(f"{len(batch)} done"), # optional progress
)
| Option | Default | Description |
|---|---|---|
batch_size |
None |
Documents per extract() call. None processes everything in one call. |
skip_processed |
True |
Passed to extract(), which applies it per batch. |
on_batch |
None |
Called with each batch's outcomes as that batch finishes. Use it for progress output or to persist partial results — a large run otherwise holds every outcome, including extracted bundles, in memory until it returns. |
The global sort is what makes batching cheap. Batching an unsorted list splits it by input order, so a later batch could carry documents older than what an earlier batch already persisted. Those are still extracted, but on the slower split-context path — an extra provenance query each, plus the reconciliation the API runs on them. Sorting first puts every batch boundary on a clean chronological cut, so each patient's documents reach the API oldest-first even when they span several batches, and none of them pay for a split a plain forward pass never needs.
Batches are cut by index, so the walk always terminates — it does not rely on
skip_processed to advance, and works with skip_processed=False too.
All documents are validated before the first batch runs, so a bad reference
late in the list surfaces before earlier batches spend anything. This is also
the only place document_id uniqueness is checked across the whole dataset:
extract() sees one batch at a time, so a duplicate pair split across two
batches would otherwise slip through.
Prefer larger batches. Each extract() call issues one FHIR query per distinct
patient in that batch, so halving batch_size roughly doubles the query
overhead. Batching bounds how much work an interruption loses; 500 bounds that
about as usefully as 50 while doing a tenth of the chatter.
extract() — a single pass
for outcome in pipeline.extract(
documents,
skip_processed=True, # default: query FHIR and skip already-processed docs
limit=10, # optional: cap this call, e.g. a sanity check before a full run
):
...
| Option | Default | Description |
|---|---|---|
skip_processed |
True |
Query FHIR for already-processed document IDs and skip them. Set to False to process all documents. |
limit |
None |
Cap the number of documents processed in this call, taken from the front of the list after skip_processed filtering. None processes everything passed. |
limit truncates a single pass; it does not chunk. extract(documents, limit=500)
on a 2000-document list processes 500 and returns — the other 1500 are
untouched. Reach for extract_all() when you want all of them.
Renamed in 0.2.0
limit was called batch_size, a name that implied chunking it never
did. Passing batch_size= to extract() raises TypeError with
migration guidance rather than being silently ignored, since ignoring it
would extract every document you passed.
When skip_processed=True, re-running extract() with the same document list is safe as long as documents have stable document_id values — reuse the identifier from your source system rather than generating a fresh one per run, or every run re-extracts everything and duplicates its resources.
Cumulative statistics
The pipeline tracks totals across all extract() calls:
| Property | Type | Description |
|---|---|---|
documents_processed |
int |
Documents successfully extracted |
documents_failed |
int |
Documents that failed extraction |
total_cost |
float |
Estimated cost in USD |
for outcome in pipeline.extract(documents):
print(outcome)
print(f"{pipeline.documents_processed} succeeded, {pipeline.documents_failed} failed")
print(f"Total cost: ${pipeline.total_cost:.3f}")
Statistics are fully updated when extract() returns. They accumulate if you call extract() multiple times on the same pipeline.
Timeouts
| Client | Default | Notes |
|---|---|---|
| Cavell API | 800s | Extraction involves many LLM calls; the API caps each call at ~300s |
| FHIR server | 30s | Per-request timeout for all FHIR operations |
Deleted patient detection
Before making any extraction API calls, extract() verifies that all referenced patients still exist on the FHIR server. If a patient was deleted (HAPI returns 404 or 410 Gone), the pipeline raises RuntimeError immediately — before spending on API calls:
RuntimeError: Patient 'MRN-12345' (FHIR id 1007) no longer exists
on the FHIR server — re-run seed() to restore it
This catches the common mistake of deleting a patient and forgetting to re-seed before extracting.
Deleting and Redoing a Patient
If a patient's extracted data looks wrong, delete all their resources and re-extract. Cascade delete removes the patient and everything referencing them — organizations and practitioners are unaffected.
# Delete
fhir_id = client.find_patient_id("MRN-12345")
client.delete_patient_resources(fhir_id)
# Re-extract: create a fresh pipeline and re-seed (recreates the deleted patient),
# then extract. skip_processed=True (default) means only the deleted patient's
# docs are re-processed.
pipeline = IngestionPipeline(client, tier="low", default_organization="CGH-001")
pipeline.seed(organizations=[...], patients=[...], practitioners=[...])
for outcome in pipeline.extract_all(all_documents, batch_size=500):
...
Important: You must create a new pipeline and re-run seed() after deleting a patient. If you try to extract with the old pipeline, extract() will detect that the patient no longer exists and raise RuntimeError.
Requires allow_cascading_deletes=true in HAPI config (set in the repo's docker-compose).
Multi-Organization
Documents from different facilities for the same patient work without extra configuration:
for outcome in pipeline.extract([
Document(
text="Initial visit at City General...",
patient_identifier="MRN-12345",
date="2024-01-15",
document_id="note-001",
organization_identifier="CGH-001",
),
Document(
text="Transfer to St. Mary's...",
patient_identifier="MRN-12345",
date="2024-02-01",
document_id="note-002",
organization_identifier="SMH-002",
),
])