An end-to-end Dagster pipeline that ingests PubMed bulk XML, filters for a target chronic condition, extracts structured treatment-evidence data using Gemini, and produces analytics stored in PostgreSQL.
- Setup Instructions
- Architecture Overview
- Directory Structure
- Condition Choice
- Schema Justification
- Prompt Design Decisions
- Analytics Design Decisions
- Known Issues and Bugs
- Suggested Improvements
- Example Queries
Prerequisites: Docker, Docker Compose, a Gemini API key (free from Google AI Studio)
git clone <your-repo-url>
cd clinical-insight-engine
cp .env.example .env
# Edit .env — set POSTGRES_PASSWORD and GEMINI_API_KEY
docker compose up --buildThe Dagster webserver starts at http://localhost:3000. The database schema is initialised automatically by init_db before the webserver boots.
-
In the Dagster UI → Automation → Sensors → toggle
pubmed_pipeline_sensorON -
Navigate to:
Jobs → pubmed_discovery_job → Launchpad
Update the runtime configuration if needed.
Default configuration:
ops:
discover_pubmed_files:
config:
no_of_files_to_load: 20
target_conditions:
- type 2 diabetesConfiguration details:
-
no_of_files_to_load
Controls how many newly discovered PubMed baseline files are processed. -
target_conditions
Filters PubMed articles by medical condition keywords during parsing and analytics.
After updating the configuration (or using the defaults):
Materialize all → Launch Run
The discovery step completes in approximately ~10 seconds.
- The sensor automatically checks every 30 seconds for newly discovered partitions and triggers:
pubmed_ingestion_and_analytics_job
The downstream pipeline automatically performs:
- PubMed file download
- XML parsing
- Study classification
- Publication trend aggregation
Running tests:
task testsThe pipeline follows a medallion architecture (Bronze → Silver → Gold) implemented as Dagster software-defined assets with dynamic partitions.
NCBI FTP Server
│
▼
[discover_pubmed_files] ← unpartitioned, registers partition keys
│ (pubmed_pipeline_sensor fires here)
▼
[download_pubmed_files] ← one partition per .xml.gz file
│ downloads + MD5 validates → bronze.pubmed_files
▼
[parse_pubmed_articles] ← one partition per file
│ streams lxml iterparse, condition-filters → silver.pubmed_articles
▼
├──► [publication_trends] → gold.publication_trends
└──► [study_classifications] → gold.article_study_classifications
| Job | Purpose |
|---|---|
pubmed_discovery_job |
Step 1 — run manually once; connects to FTP and registers new filenames as dynamic partition keys |
pubmed_ingestion_and_analytics_job |
Step 2 — auto-launched by sensor; runs download → parse → analytics for each partition |
analytical_pubmed_job |
Targeted re-run of gold layer only (useful when silver is already populated) |
Dagster requires partition keys to exist before a job can be scheduled against them. discover_pubmed_files creates those keys, so it cannot be in the same job as the partitioned assets — doing so causes the UI to present a partition selection dialog before any keys exist. The sensor bridges the gap: it watches discover_pubmed_files for a successful materialization, reads the newly registered keys from output metadata, and emits one RunRequest per key.
clinical-insight-engine/
├── Dockerfile # python:3.12-slim + uv sync --frozen
├── docker-compose.yml # db + dagster_webserver + dagster_daemon
├── pyproject.toml # single source of truth for all dependencies
├── uv.lock # fully pinned lockfile — reproducible installs
├── workspace.yaml # Dagster workspace pointing at clinical_insight.definitions
├── dagster.yaml # Dagster instance config (postgres backend, log retention)
├── init.sql # creates bronze/silver/gold postgres schemas on first boot
├── .env.example # all required env vars documented
├── storage/ # bind-mounted volume; .xml.gz files downloaded here
│ └── downloads/
│
├── clinical_insight/ # main Python package
│ ├── definitions.py # Dagster Definitions — wires assets, jobs, sensors, resources
│ │
│ ├── apps/
│ │ ├── resources.py # ConfigurableResource subclasses (Database, Gemini, Storage)
│ │ └── partitions.py # DynamicPartitionsDefinition for raw PubMed files
│ │
│ ├── core/
│ │ ├── config.py # pydantic-settings Settings (reads .env)
│ │ ├── pubmed_config.py # Dagster Config: target_condition, max_articles_per_file
│ │ └── logging.py # structured logging setup
│ │
│ ├── clients/
│ │ └── pubmed/
│ │ ├── ftp_client.py # FTP connection, file listing, download (with tenacity retry)
│ │ ├── parser.py # lxml iterparse streaming XML parser
│ │ ├── extractors.py # field-level extraction from PubmedArticle XML elements
│ │ └── checksum.py # MD5 verification helper
│ │
│ ├── db/
│ │ ├── base.py # SQLAlchemy declarative base
│ │ ├── session.py # engine construction from settings
│ │ ├── init_db.py # creates schemas + tables (called before webserver starts)
│ │ └── models/
│ │ ├── bronze.py # PubmedFile — download registry and status tracking
│ │ ├── silver.py # PubmedArticle, ArticleAuthor, ArticleGrant
│ │ ├── gold.py # PublicationTrend, ArticleStudyClassification
│ │ └── analytics.py # (reserved for future analytics models)
│ │
│ ├── assets/
│ │ ├── bronze/
│ │ │ ├── discover_pubmed_files.py # FTP listing → dynamic partition registration
│ │ │ └── download_pubmed_file.py # per-partition download + MD5 + bronze record
│ │ ├── silver/
│ │ │ └── parse_pubmed_articles.py # per-partition XML parse + condition filter + persist
│ │ └── gold/
│ │ ├── publication_trends.py # metadata-powered: year-over-year publication counts
│ │ └── study_classifications.py # LLM-powered: async Gemini study type classification
│ │
│ ├── services/
│ │ ├── ingestion/
│ │ │ ├── download_service.py # orchestrates FTP download + validation
│ │ │ ├── file_registry_service.py # bronze.pubmed_files CRUD
│ │ │ ├── pubmed_article_ingestion_service.py # parse → filter → persist loop
│ │ │ └── condition_filter_service.py # text normalisation + condition substring match
│ │ ├── analytics/
│ │ │ └── publication_trend_service.py # full-refresh aggregation into gold.publication_trends
│ │ └── llm/
│ │ ├── schemas.py # StudyClassificationResult Pydantic model
│ │ └── study_classification_service.py # prompt construction + response validation
│ │
│ └── jobs/
│ └── pubmed_jobs.py # job and sensor definitions (see architecture note above)
│
└── tests/
├── conftest.py # shared fixtures: sample XML, mock DB session
├── test_ingestion.py # XML parse + asset wrapper stubs
├── test_llm_extraction.py # LLM mock + extraction shape tests
└── test_db.py # model validation stubs
clients/ vs services/ — clients own the protocol boundary (FTP, XML, HTTP). Services compose clients into business logic (download-and-validate, ingest-and-filter). Assets orchestrate services and handle Dagster metadata. This separation means the FTP client and XML parser are independently testable without Dagster context.
apps/ — Dagster-specific configuration (resources, partition definitions) lives here to keep the clinical_insight core importable in tests without pulling in Dagster machinery.
core/ — settings and logging are isolated so any module can import them without circular dependencies.
db/models/ split by layer — bronze, silver, and gold models in separate files makes the medallion contract explicit and prevents accidental cross-layer imports.
Target condition: Type 2 Diabetes
Type 2 diabetes was chosen for three practical reasons:
-
Volume. It is one of the most-published chronic conditions in PubMed. Two baseline files (~60,000 total records) are expected to yield well over 2,000 matching abstracts after filtering, comfortably within the 2,000–5,000 target range without needing many files.
-
Treatment diversity. The literature spans metformin (decades of evidence), GLP-1 agonists, SGLT2 inhibitors, insulin regimens, and lifestyle interventions — a rich treatment landscape that makes the analytics genuinely interesting.
-
Clear MeSH coverage. "Diabetes Mellitus, Type 2" is a well-established MeSH descriptor with consistent historical usage, making the keyword/MeSH filter reliable across publication years.
Filtering approach: ConditionFilterService normalises the target condition string and the concatenated MeSH terms + title + abstract to lowercase alphanumeric tokens, then checks for a substring match. This is intentionally simple: it catches "type 2 diabetes", "diabetes mellitus type 2", and MeSH term variants without requiring a curated synonym list. The trade-off is occasional false positives (articles that mention T2D as a comorbidity rather than the primary focus) and false negatives (articles using abbreviations like "T2DM" without the full term). A production system would use PubMed's own MeSH qualifier filtering via the Entrez API or a pre-built synonym map.
Files needed: The pipeline defaults to 2 baseline files per discovery run (hardcoded safety cap in discover_pubmed_files). Each file contains ~30,000 records; with a ~5–8% T2D hit rate, 2 files produces roughly 3,000–5,000 matched abstracts. Remove the new_files = new_files[:2] slice to ingest the full baseline.
The database uses three PostgreSQL schemas mirroring the medallion layers. The init.sql creates them before SQLModel creates the tables.
Tracks every baseline .xml.gz file downloaded from NCBI. The download_status / parse_status enum columns allow the pipeline to resume after partial failures — a file that downloaded successfully but failed to parse retains DOWNLOADED status so parse can be retried without re-downloading. retry_count and error_message are reserved for future retry logic. The filename column is the natural idempotency key: uniquely constrained so re-running a partition is always safe.
One row per article that passed the condition filter. pmid is uniquely constrained — the primary deduplication guard. source_file_id preserves provenance back to the bronze record. mesh_terms, keywords, and publication_types are stored as PostgreSQL ARRAY(String) rather than normalised junction tables: for read-heavy analytics (co-occurrence, filtering) array columns with GIN indexes outperform joins, and PubMed's term lists are small enough (typically <30 entries) that denormalisation is safe. condition_matched stores the exact condition string used at ingest time so the gold layer can filter without joining back to bronze.
ArticleAuthor and ArticleGrant are separate tables with foreign keys to pubmed_articles. Author position is stored so first/last author analyses are possible. The UniqueConstraint("article_id", "author_position") prevents duplicate author rows on re-ingestion.
A pre-aggregated summary table: one row per (condition_name, publication_year). The analytics asset runs a full delete-then-insert on every materialisation so the counts always reflect the current silver state. The UniqueConstraint enforces this invariant at the database level.
One row per classified article with the LLM's study type, confidence score, and model name. Storing the model name enables future comparisons across model versions. The UniqueConstraint("article_id") ensures one classification per article even if the asset is re-run.
The study classification prompt (study_classification_service.py) was designed around three goals:
1. Closed taxonomy with explicit examples
The prompt enumerates exactly 10 allowed study types and provides three few-shot examples that cover the most common ambiguities (RCT vs Clinical Trial, systematic review vs narrative review, animal model). Without a closed list, Gemini produces varied phrasings ("randomised trial", "RCT", "Randomised Controlled Trial") that make downstream grouping unreliable. With it, the model's output maps directly to enum values.
2. Confidence score as a calibration signal
Asking for a confidence_score alongside study_type serves two purposes: it gives the pipeline a way to flag low-confidence classifications for human review, and it pressures the model to be more deliberate — producing "Unknown" at 0.4 confidence rather than a plausible-but-wrong label at 0.9.
3. JSON-only output with fence stripping
temperature=0.0 is used for determinism. The response is stripped of markdown code fences (```json … ```) before parsing because Gemini frequently wraps JSON in fences even when the prompt says not to. The fallback validates that study_type falls within ALLOWED_STUDY_TYPES and clamps confidence_score to [0.0, 1.0] — both guards needed in practice because the model occasionally hallucinates values outside these ranges.
Limitations acknowledged:
- The condition filter uses only the title + abstract + MeSH terms. Articles with results presented only in figures or tables will not surface any treatment names to the LLM.
- Study type classification on short or non-English abstracts degrades. Articles without an abstract receive the string "No abstract available" which reliably produces "Unknown".
- The prompt does not include the full PubMed publication type field (e.g.,
[Publication Type = Randomized Controlled Trial]), which is a stronger signal than abstract text. Cross-referencing LLM output with publication type would improve precision significantly.
What it does: Counts articles per publication year for the target condition and writes to gold.publication_trends. The full-refresh approach means the gold table always reflects the current silver state — no incremental merge logic needed.
Limitations: Year is extracted from PubDate/Year. When PubMed only provides a MedlineDate string (e.g., "2003 Jul-Aug"), the extractor parses the first 4-digit token. Articles with no parseable year are excluded from trend counts.
What it does: For each unclassified silver article, sends title + abstract to Gemini and persists the study type + confidence score. Requests are batched concurrently using asyncio.Semaphore (default concurrency: 5) to reduce wall-clock time without overwhelming the API.
Batch limit: The asset processes at most 50 unclassified articles per materialisation. This cap prevents timeouts on large silver tables and makes re-runs incremental. Remove the .limit(batch_limit) call for a full corpus run.
Limitations: The classification is at the abstract level, not the article level. A meta-analysis abstract that discusses RCT methodology may be classified as "Randomized Controlled Trial". Confidence scores are model-generated and not calibrated against ground truth; they should be treated as relative signals rather than probabilities.
Run these against the clinical_insight database after a successful pipeline run:
-- Publication volume by year for the target condition
SELECT
publication_year,
publication_count
FROM gold.publication_trends
WHERE 'type 2 diabetes' = ANY(condition_names)
ORDER BY publication_year;
-- Study type distribution
SELECT
c.study_type,
COUNT(*) AS count,
ROUND(
AVG(c.confidence_score)::numeric,
3
) AS avg_confidence
FROM gold.article_study_classifications c
JOIN silver.pubmed_articles a
ON a.id = c.article_id
WHERE 'type 2 diabetes' = ANY(a.condition_matched)
GROUP BY c.study_type
ORDER BY count DESC;
-- Top funding agencies
SELECT
g.agency,
COUNT(*) AS grant_count
FROM silver.article_grants g
JOIN silver.pubmed_articles a
ON a.id = g.article_id
WHERE g.agency IS NOT NULL
AND 'type 2 diabetes' = ANY(a.condition_matched)
GROUP BY g.agency
ORDER BY grant_count DESC
LIMIT 20;
-- Most common MeSH terms across matched articles
SELECT
UNNEST(mesh_terms) AS term,
COUNT(*) AS frequency
FROM silver.pubmed_articles
WHERE 'type 2 diabetes' = ANY(condition_matched)
GROUP BY term
ORDER BY frequency DESC
LIMIT 30;
-- Articles by publication year with study classification
SELECT
a.publication_year,
c.study_type,
COUNT(*) AS count
FROM silver.pubmed_articles a
JOIN gold.article_study_classifications c
ON c.article_id = a.id
WHERE 'type 2 diabetes' = ANY(a.condition_matched)
GROUP BY
a.publication_year,
c.study_type
ORDER BY
a.publication_year,
count DESC;
-- Research geography: top countries by author affiliation (requires text extraction)
SELECT
country,
COUNT(*) AS article_count
FROM silver.pubmed_articles
WHERE country IS NOT NULL
AND 'type 2 diabetes' = ANY(condition_matched)
GROUP BY country
ORDER BY article_count DESC;
-- Files ingested and their parse status
SELECT filename, download_status, parse_status, parsed_at
FROM bronze.pubmed_files
ORDER BY discovered_at;