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.
Cadence · 01
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.
The seam · 02
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.
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.
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.
RevFind-Daily-Processing
-
▶RevFind-Source-Load-PM-Data— fan-out over eCW PM/EDI export files
- ExecuteJobForFiles per file mask:
*ExportChargeItems*,*ExportEdiInvEob*,*ExportEdiInvCPT*,*ExportEdiInvDiagnosis*,*ExportDoctors*, …
- ExecuteJobForFiles per file mask:
-
jobsTransform-InsuranceMap-ECW · Transform-ChargeType · Transform-AccountsShape the eCW export into the RevFind account/charge model
-
shared tailProcess-Entities → DimensionMembers → Accounts → Post-Load-Routing → Scoring → ContractModeling → ReportsThe common RevFind pipeline: score accounts vs contracts, model, report → shared with ModMed
Flow-ECW-Processing
-
▶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.
-
jobsLoadEntities · Transform · LoadPatients · LoadDimensionMembers · LoadVisitsBuild visits, benefits, letters (the shared ClarityFlow core)
-
workqueue + send_hl7WorkqueueAssignment · SendHL7Letters
RevFind-Daily-Processing-ModMed
-
upstreamSource*MM tablesPopulated by the modern gGastro sync (below), not a legacy file load
-
jobsTransform-InsuranceMap-ModMed · ChargeType-ModMed · Transform-Accounts-ModMed · Process-Entities-ModMedModMed-specific transforms over
Source*MM -
shared tailDimensionMembers → Accounts → Routing → Scoring → ContractModeling → ReportsSame shared RevFind tail as the eCW chain
Flow-MM-Processing
-
jobsLoadEntities · Transform · LoadPatients · LoadDimensionMembers · LoadVisitsStarts at
LoadEntities— noSource-LoadDatastep. The source tables are already filled by the medallion sync. -
workqueue + send_hl7WorkqueueAssignment · SendHL7Letters
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.
onegi-ingest → bronze
-
dltExtract gGastro MSSQL → S3 Parquet6 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 intables.yaml. -
146 asset_checksBronze data-quality gateType-drift (
no_variant_columns), clamps, money/uuid normalization
onegi-transform → silver
-
DuckDBConform bronze → silver entitiesPer-
db_keyDuckDB 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.
onegi-transform → sync_to_source
-
BCP + sp_renameAtomic swap into legacy source tablesGlob + dedupe silver → BCP into
[dbo].[<target>__New]→sp_renameswap 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
Reference · Notes
__New → sp_rename swap. Feeds the legacy source tables.[dbo].[Source…MM] tables. Named "MM" (ModMed), now populated by gGastro via the medallion sync.tn002). One bronze/silver partition + one RunRequest per daily tick.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.