Delta Lake Export
EAP

 

The Delta Lake Export module extends Realtime Export to write FHIR resource changes to Delta Lake tables. This enables customers on data lake architectures (Databricks, Spark, cloud object stores) to receive FHIR data directly, without building custom ETL pipelines.

The module connects to a customer-managed Spark Connect server over gRPC and writes using version-guarded MERGE / UPDATE / DELETE: the table holds one row per resource, replaced in place on each update and physically removed on delete. Consumers query the table directly with a plain SELECT.

Architecture
EAP

 

Delta Lake Export is configured as a separate module type (REALTIME_EXPORT_DELTA_LAKE) with its own configuration. It shares the same rules model as the JDBC-based export, but writes to Delta Lake instead of a database. Each module owns its own producer interceptor and publishes to a module-specific broker channel, so the two can operate independently with no shared state or coupled configuration.

The module connects to an external Spark Connect server (the customer's Spark cluster) via gRPC.

FHIR Storage Module
        |
        v
  Interceptor (one per export module, CREATE/UPDATE/DELETE hooks)
   |                          |
   v                          v
 Channel                    Channel
 "realtime.export"          "realtime.export.deltalake"
   |                          |
   v                          v
 JDBC Realtime Export       Delta Lake Export
                              |
                              v
                        Spark Connect (gRPC)
                              |
                              v
                   Customer-managed Spark cluster
                    (Docker, k8s, Databricks, EMR)
                              |
                              v
                      Delta Lake Tables
                     (Parquet + _delta_log)

Deployment

The customer provides a Spark Connect server. Options include:

  • Docker Compose — apache/spark:4.0.0 with the delta-connect-server plugin (see example compose file in the test resources)
  • Kubernetes sidecar — a Spark Connect pod alongside CDR
  • Managed service — Databricks Connect, Amazon EMR, or Google Dataproc

CDR connects via the deltalake.spark_connect_url configuration property (e.g. sc://spark-host:15002).

Write Semantics
EAP

 

Each FHIR resource operation is written with a version-guarded Delta MERGE:

  • Create / Update: a version-guarded MERGE statement (shown below). Only updates whose source version is strictly greater than the target's replace the existing row, so out-of-order delivery is tolerated — stale messages are dropped.
  • Delete: parameterized SQL DELETE on the row's id. Deletes on non-existent tables are no-ops.
MERGE INTO target t USING source s ON t.id = s.id
WHEN MATCHED AND t.version < s.version THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *

Every parent row carries two system columns:

ColumnTypeDescription
idSTRINGUnqualified, versionless resource ID (e.g., Patient/123)
versionINTFHIR resource version

Consumers can query current state with a plain SELECT:

SELECT * FROM patient_table WHERE id = 'Patient/123'

The trade-off is no built-in history on the primary table; consumers that need an audit trail should enable retainAllHistory (separate _history tables) or capture history upstream.

Configuration
EAP

 

The Delta Lake Export module reuses the same rules JSON format as the JDBC Realtime Export module. See Realtime Export Rules Definition for details on configuring transformers, columns, and FHIRPath expressions.

Module Properties

module.realtime_export_delta_lake.config.deltalake.path=/path/to/delta/tables
module.realtime_export_delta_lake.config.deltalake.spark_connect_url=sc://spark-host:15002
module.realtime_export_delta_lake.config.script.text={"transformers": [...]}
module.realtime_export_delta_lake.config.script.file=/path/to/rules.json
module.realtime_export_delta_lake.config.channel.concurrent_consumers=1
module.realtime_export_delta_lake.config.channel.concurrent_retry_consumers=1
module.realtime_export_delta_lake.config.channel.prefix=
PropertyDescriptionDefault
deltalake.pathPath where Delta tables are created (local filesystem, s3a://, abfss://, gs:// — depends on Spark cluster configuration)(required)
deltalake.spark_connect_urlSpark Connect server URL (e.g. sc://spark-host:15002)(required)
script.textInline JSON rules configurationEmpty rules
script.filePath to external JSON rules file(none)
channel.concurrent_consumersNumber of concurrent message consumers1
channel.concurrent_retry_consumersNumber of concurrent retry consumers1
channel.prefixChannel name prefix(empty)

Table Creation

Tables are created automatically on first write. The table schema is derived from the rules configuration, with system columns added automatically. New tables are created with auto-optimize TBLPROPERTIES enabled (delta.autoOptimize.optimizeWrite=true, delta.autoOptimize.autoCompact=true) so write batches are compacted automatically.

Child Tables
EAP

 

Repeating FHIR elements (e.g., Patient.name, Patient.address) are configured as child tables in the rules JSON. The Delta Lake export module supports two modes for handling child tables, controlled by the flatten property on the child table configuration:

flattenBehavior
trueThe child elements are embedded in the parent row as a nested ARRAY&lt;STRUCT&lt;...&gt;&gt; Parquet column. Single Delta table per resource.
false (default)Each child element is written as a row in a separate Delta table with link columns.

Flattened Child Tables (flatten: true)

When flatten is true, the child table's tableName becomes the name of a nested array column in the parent table. The child transformer's columns become the fields of the struct inside the array.

{
  "transformers": [
    {
      "resourceType": "Patient",
      "tableName": "patient_table",
      "columns": [
        {"columnName": "is_active", "fhirPath": "active", "columnType": "BOOLEAN"}
      ],
      "childTables": [
        {
          "fhirPath": "Patient.name",
          "tableName": "names",
          "flatten": true,
          "childTransformer": {
            "columns": [
              {"columnName": "family", "fhirPath": "family", "columnType": "STRING"},
              {"columnName": "use", "fhirPath": "use", "columnType": "STRING"}
            ]
          }
        }
      ]
    }
  ]
}

The resulting patient_table schema includes a names column of type ARRAY&lt;STRUCT&lt;family: STRING, use: STRING&gt;&gt;. Spark consumers can query it via EXPLODE:

SELECT id, n.family, n.use
FROM patient_table
LATERAL VIEW EXPLODE(names) AS n
WHERE n.use = 'official'

Recursive nesting is supported — child tables can have their own flattened child tables.

Separate Child Tables (flatten: false)

When flatten is false (the default), each child element becomes a row in its own Delta table. These child tables include three link columns:

ColumnDescription
idAuto-generated UUID for the child row
parent_referenceThe ID of the parent row this child belongs to
source_resource_idThe unqualified ID of the top-level FHIR resource

Non-flattened child tables use replace-children semantics: on every parent create or update, all existing child rows for the resource are deleted, then the new child rows are inserted. On parent delete, child rows are cascade-deleted before the parent row is removed.

Schema Evolution
EAP

 

Schema evolution is automatic. When the rules configuration adds a new column that is not present in the live Delta table, the writer runs ALTER TABLE ... ADD COLUMNS automatically. Existing rows get null for the new column.

Schema changeBehavior
New column added to rulesAutomatic ALTER TABLE ADD COLUMNS. Old rows get null.
Column removed from rulesColumn persists in the table; new rows write null. No action needed.
Type change on an existing columnFails fast with a clear error naming the column and both types. Operator must resolve manually.

Schema comparison happens once on the first write to each table, then the resolved schema is cached for the lifetime of the module instance.

Type Mapping
EAP

 

The following table shows how RTE column types map to Delta Lake (Parquet) types:

RTE Column TypeDelta Lake TypeNotes
STRINGStringType
INTIntegerType
LONGLongType
BOOLEANBooleanType
FLOATFloatType
DOUBLEDoubleType
DATE_ONLYDateTypeStored as epoch day
DATE_TIMESTAMPTimestampTypeStored as microseconds since epoch
BLOBBinaryType
CLOBStringTypeMapped to string in Delta Lake

Differences from JDBC Realtime Export
EAP

 
AspectJDBC ExportDelta Lake
Write patternINSERT/UPDATE/DELETE (mutable)MERGE / UPDATE / DELETE (mutable, current state)
Current stateAlways up to date in tableAlways up to date in table
HistoryOptional (retainAllHistory)Optional (retainAllHistory) — separate _history tables
History tablesSeparate _history tablesSeparate _history tables (when retainAllHistory: true)
Transaction atomicityDatabase transactionPer-row only
Schema creationManual (pre-create tables)Automatic on first write
Schema evolution (add column)Manual ALTER TABLEAutomatic ALTER TABLE ADD COLUMNS
Non-flattened child tablesSupportedSupported (replace-children)
Process footprint in CDRJDBC driver onlySpark Connect client (~33 MB)
External dependencyDatabaseSpark Connect server
TargetJDBC databasesDelta Lake (Parquet + transaction log)

Limitations
EAP

 
  • No batching: Each resource change is written as a separate Spark Connect operation. Write batching to coalesce multiple changes into a single Delta commit is deferred and will be provided in a future release.
  • No partitioning: Delta tables are not partitioned. Partitioning strategy is planned for a future release.
  • Spark Connect server required: CDR connects to an external Spark Connect server via gRPC. The customer must provision and manage this server. See Deployment for options.