← Pipeline explainers OneGI · gastroenterology group

MDClarity · Pipeline explainer

How OneGI's scheduled pipelines work

OneGI is where the legacy → medallion migration seam is on full display. It runs both MDClarity products (RevFind + ClarityFlow) across three source systems on two stacks: legacy file/DB pipelines for eCW and ModMed on WORKER-1, and a modern Dagster medallion pipeline for gGastro that reverse-ETLs into the very tables the legacy ModMed processing reads.

Legacy: WORKER-1 Task Scheduler → RunDriverJob.ps1 → jobs on prod-db-3
Modern: Dagster workspace ~/src/mdc/MdClarity/pipelines (onegi-ingest + onegi-transform) → S3 datalake + OneGIDestination MSSQL
Sources: legacy export (dev-1 2026-07-09) + the modern dg workspace · Structure & names only, no patient data

Sources
3
eCW · ModMed · gGastro
Products
2
RevFind + ClarityFlow, per source
Legacy daily tasks
2
ONEGI @ 23:00 · ONEGI-MM @ 03:00
gGastro source DBs
6
per-practice MSSQL → medallion
Flow / ClarityFlow
RevFind
modern medallion (Dagster)
legacy↔modern seam
▶ click sub-jobs to expand

Cadence · 01

When each pipeline runs

The legacy tasks are simple once-daily runs on WORKER-1 (no repetition). The modern medallion side runs its own daily bronze pull, then chains silver + sync event-driven off bronze success — so gGastro data is landed and pushed into the legacy tables before the ModMed task reads them.

00:00 06:00 12:00 18:00 24:00
gGastro bronze
modern · daily
06:00
silver → sync
modern · event
on bronze success →
ModMed (MM)
legacy · daily
03:00
eCW
legacy · daily
23:00

The seam · 02

Two stacks, one destination

OneGI is mid-migration. The eCW path is fully legacy (files → WORKER-1). The ModMed (MM) path is a hybrid: its downstream processing is still the legacy MDClarity pipeline, but its source data no longer loads from files — it's produced by the modern gGastro medallion stack and pushed into the same SQL Server Source*MM tables the legacy jobs read.

⇄ legacy ↔ modern bridge
The modern onegi-transform sync (reverse-ETL) BCPs each silver entity into [dbo].[<target>__New] and atomically sp_rename-swaps it over the live [dbo].[Source*MM]. The legacy Flow-MM-Processing then starts at LoadEntities — it has no source-load step at all (unlike Flow-ECW-Processing, which still loads eCW files). That missing step is the tell: gGastro on the medallion stack has already filled the source tables.
Legacy stack · WORKER-1 · prod-db-3

Two daily Windows tasks, each running its RevFind chain then its Flow chain. The two RevFind chains converge on a shared scoring / contract-modeling / reporting tail.

ONEGI · 23:00 · eCW

RevFind-Daily-Processing

Action 1 · "Process ECW Accounts"
11-job RevFind chain
  • ▶RevFind-Source-Load-PM-Data— fan-out over eCW PM/EDI export files
    • ExecuteJobForFiles per file mask: *ExportChargeItems*, *ExportEdiInvEob*, *ExportEdiInvCPT*, *ExportEdiInvDiagnosis*, *ExportDoctors*, …
  • jobsTransform-InsuranceMap-ECW · Transform-ChargeType · Transform-Accounts
    Shape the eCW export into the RevFind account/charge model
  • shared tailProcess-Entities → DimensionMembers → Accounts → Post-Load-Routing → Scoring → ContractModeling → Reports
    The common RevFind pipeline: score accounts vs contracts, model, report → shared with ModMed
ONEGI · 23:00 · eCW

Flow-ECW-Processing

Action 2 · ClarityFlow HL7/visit pipeline
has legacy source-load
  • ▶Flow-ECW-Source-LoadData— 15× ExecuteJobForFiles + merge, with pre-parse file cleaners
    • One sub-job per eCW clinical file type (ProblemList, Referral, PriorDiagnosis, FlowPatients, …)
    • fixFiles.py (via pyFixFiles.bat) — re-stitches ^|^ rows split by embedded newlines in free-text fields back to a fixed column count, then caps the notes column at 3999 chars. Parameterized per file: ProblemList (20 cols, notes@3), Referral (9 cols, notes@5).
    • fixPriorDiagnosis.py (via pyPriorDiagnosis.bat) — the same row-reassembly, hardcoded for PriorDiagnosis (7 cols, notes@2, 3999-char cap).
    • CleanFiles.ps1 — FlowPatients: strips embedded quotes + NUL bytes + trailing blank lines.
    eCW's pipe-delimited exports carry raw newlines inside clinical free-text, stray quotes, and NULs that break delimited parsing — these cleaners normalize each file before load. The 3999-char notes cap matches the target SQL column width.
  • jobsLoadEntities · Transform · LoadPatients · LoadDimensionMembers · LoadVisits
    Build visits, benefits, letters (the shared ClarityFlow core)
  • workqueue + send_hl7WorkqueueAssignment · SendHL7Letters
ONEGI-MM · 03:00 · ModMed

RevFind-Daily-Processing-ModMed

Action 1 · "Process MM Accounts"
source via medallion ⇄
  • upstreamSource*MM tables
    Populated by the modern gGastro sync (below), not a legacy file load
  • jobsTransform-InsuranceMap-ModMed · ChargeType-ModMed · Transform-Accounts-ModMed · Process-Entities-ModMed
    ModMed-specific transforms over Source*MM
  • shared tailDimensionMembers → Accounts → Routing → Scoring → ContractModeling → Reports
    Same shared RevFind tail as the eCW chain
ONEGI-MM · 03:00 · ModMed

Flow-MM-Processing

Action 2 · ClarityFlow HL7/visit pipeline
NO source-load ⇄
  • jobsLoadEntities · Transform · LoadPatients · LoadDimensionMembers · LoadVisits
    Starts at LoadEntities — no Source-LoadData step. The source tables are already filled by the medallion sync.
  • workqueue + send_hl7WorkqueueAssignment · SendHL7Letters
Modern medallion stack · Dagster · gGastro

A Dagster workspace (onegi-ingest + onegi-transform, shared onegi-common) that lands gGastro into an S3 datalake, conforms it in DuckDB, and reverse-ETLs it into the legacy SQL Server source tables. gGastro is sharded across 6 per-practice MSSQL databases (tn002, tn031, va007, in010, oh034, ms004) — one partition each.

Bronze · 06:00 daily

onegi-ingest → bronze

dlt (connectorx) · bronze_daily_job · daily_bronze_schedule 0 6 * * *
73 tables · 146 checks
  • dltExtract gGastro MSSQL → S3 Parquet
    6 source DBs → bronze/ggastro/db={db_key}/tables/dt=YYYY-MM-DD/
    Split into 4 parallel @dlt_assets (lookups, patient, billing_charges, billing_transactions) so a memory blow-up on the biggest billing table can't take the rest of bronze down. Per-table column projections in tables.yaml.
  • 146 asset_checksBronze data-quality gate
    Type-drift (no_variant_columns), clamps, money/uuid normalization
Silver · event-driven

onegi-transform → silver

DuckDB · silver_job · sensor bronze_daily_succeeded
18 entities · per-db_key
  • DuckDBConform bronze → silver entities
    Per-db_key DuckDB SQL over bronze Parquet → silver/ggastro/db={db_key}/{entity}/ (SCD Type-1 merge)
    18 entities: account · adjustment · carrier · charge · claim · location · modifier · organization · patient · patient_insurance · payer · payment · provider · refund · transaction · usa_address — plus GI-specific colonoscopy_appointment_classification & appointment_reschedule_chain. Checks: PK uniqueness, composite keys, SCD null handling.
Sync → source · event-driven

onegi-transform → sync_to_source

reverse-ETL · sync_to_source_job · AutomationCondition.eager()
→ Source*MM ⇄
  • BCP + sp_renameAtomic swap into legacy source tables
    Glob + dedupe silver → BCP into [dbo].[<target>__New] → sp_rename swap over live [dbo].[Source*MM]
    Sch-M lock for milliseconds; idempotent, so duplicate firings during burst windows are wasted compute, never data damage. This is the hand-off point to the legacy ModMed processing.

Reference · Interfaces

Who OneGI talks to
EndpointStackDirectionCarries
mdcssh SFTPlegacyinbound ⬇ eCW PM/EDI + clinical export files
gGastro — 6 MSSQL DBsmoderninbound ⬇ dlt extract (tn002/tn031/va007/in010/oh034/ms004)
S3 datalake customer/OneGImodernbronze/silver Parquet: bronze/ggastro/…, silver/ggastro/…
OneGIDestination MSSQL [dbo].[Source*MM]seam ⇄reverse-ETL ⬆ Conformed gGastro → the tables legacy ModMed reads
Patientslegacyoutbound ⬆ ClarityFlow estimate letters (SendHL7Letters)

Reference · Notes

Terms & open items
bronze / silver
Medallion layers. Bronze = raw source landed as Parquet in S3; silver = conformed, deduped entities (SCD Type-1) in DuckDB.
dlt
The ingest library (connectorx backend) that replicates gGastro MSSQL tables to bronze Parquet, partitioned by source DB.
reverse-ETL / sync_to_source
Pushes silver back into SQL Server via BCP → __New → sp_rename swap. Feeds the legacy source tables.
Source*MM
The [dbo].[Source…MM] tables. Named "MM" (ModMed), now populated by gGastro via the medallion sync.
db_key
A gGastro practice database (e.g. tn002). One bronze/silver partition + one RunRequest per daily tick.
shared RevFind tail
Process-Entities → DimensionMembers → Accounts → Post-Load-Routing → Scoring → ContractModeling → Reports; shared by the eCW and ModMed RevFind chains.
eCW cleaner scripts
Pre-parse cleaners under D:\ClientData\OneGI\ECW\: fixFiles.py / fixPriorDiagnosis.py re-stitch rows split by newlines in free-text and cap notes at 3999 chars; CleanFiles.ps1 strips quotes / NULs. eCW exports are newline- and quote-dirty.
gold / serving
Not in this repo's scope — downstream of the SQL Server destination; the legacy RevFind reporting is the current consumer.