using Parquet; using Parquet.Data; using Parquet.Schema; public static class ParquetFileBuilder { public static async Task BuildObservationsAsync(IReadOnlyList rows) { var schema = new ParquetSchema( new DataField("observation_id"), new DataField("encounter_id"), new DataField("patient_id"), new DataField("mrn"), new DataField("observation_code"), new DataField("value"), new DataField("unit"), new DataField("source"), new DataField("recorded_at"), new DataField("kafka_partition"), new DataField("kafka_offset") ); using var ms = new MemoryStream(); using (var writer = await ParquetWriter.CreateAsync(schema, ms)) using (var rg = writer.CreateRowGroup()) { var f = schema.DataFields; await rg.WriteColumnAsync(new DataColumn(f[0], rows.Select(r => r.ObservationId).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[1], rows.Select(r => r.EncounterId).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[2], rows.Select(r => r.PatientId).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[3], rows.Select(r => r.Mrn).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[4], rows.Select(r => r.ObservationCode).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[5], rows.Select(r => r.Value).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[6], rows.Select(r => r.Unit).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[7], rows.Select(r => r.Source).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[8], rows.Select(r => r.RecordedAt).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[9], rows.Select(r => r.KafkaPartition).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[10], rows.Select(r => r.KafkaOffset).ToArray())); } return ms.ToArray(); } public static async Task BuildAlertsAsync(IReadOnlyList rows) { var schema = new ParquetSchema( new DataField("alert_id"), new DataField("encounter_id"), new DataField("patient_id"), new DataField("alert_type"), new DataField("severity"), new DataField("details"), new DataField("triggered_at"), new DataField("kafka_partition"), new DataField("kafka_offset") ); using var ms = new MemoryStream(); using (var writer = await ParquetWriter.CreateAsync(schema, ms)) using (var rg = writer.CreateRowGroup()) { var f = schema.DataFields; await rg.WriteColumnAsync(new DataColumn(f[0], rows.Select(r => r.AlertId).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[1], rows.Select(r => r.EncounterId).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[2], rows.Select(r => r.PatientId).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[3], rows.Select(r => r.AlertType).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[4], rows.Select(r => r.Severity).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[5], rows.Select(r => r.Details).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[6], rows.Select(r => r.TriggeredAt).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[7], rows.Select(r => r.KafkaPartition).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[8], rows.Select(r => r.KafkaOffset).ToArray())); } return ms.ToArray(); } public static async Task BuildEncountersAsync(IReadOnlyList rows) { var schema = new ParquetSchema( new DataField("encounter_id"), new DataField("patient_id"), new DataField("previous_status"), new DataField("new_status"), new DataField("changed_at"), new DataField("kafka_partition"), new DataField("kafka_offset") ); using var ms = new MemoryStream(); using (var writer = await ParquetWriter.CreateAsync(schema, ms)) using (var rg = writer.CreateRowGroup()) { var f = schema.DataFields; await rg.WriteColumnAsync(new DataColumn(f[0], rows.Select(r => r.EncounterId).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[1], rows.Select(r => r.PatientId).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[2], rows.Select(r => r.PreviousStatus).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[3], rows.Select(r => r.NewStatus).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[4], rows.Select(r => r.ChangedAt).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[5], rows.Select(r => r.KafkaPartition).ToArray())); await rg.WriteColumnAsync(new DataColumn(f[6], rows.Select(r => r.KafkaOffset).ToArray())); } return ms.ToArray(); } }