Files
vigilcare-clinical/VigilCareClinicalAPI/DataLake/ParquetFileBuilder.cs
T

102 lines
5.3 KiB
C#

using Parquet;
using Parquet.Data;
using Parquet.Schema;
public static class ParquetFileBuilder
{
public static async Task<byte[]> BuildObservationsAsync(IReadOnlyList<ObservationRow> rows)
{
var schema = new ParquetSchema(
new DataField<string>("observation_id"),
new DataField<string>("encounter_id"),
new DataField<string>("patient_id"),
new DataField<string>("mrn"),
new DataField<string>("observation_code"),
new DataField<double>("value"),
new DataField<string>("unit"),
new DataField<string>("source"),
new DataField<string>("recorded_at"),
new DataField<int>("kafka_partition"),
new DataField<long>("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<byte[]> BuildAlertsAsync(IReadOnlyList<AlertRow> rows)
{
var schema = new ParquetSchema(
new DataField<string>("alert_id"),
new DataField<string>("encounter_id"),
new DataField<string>("patient_id"),
new DataField<string>("alert_type"),
new DataField<string>("severity"),
new DataField<string>("details"),
new DataField<string>("triggered_at"),
new DataField<int>("kafka_partition"),
new DataField<long>("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<byte[]> BuildEncountersAsync(IReadOnlyList<EncounterStatusRow> rows)
{
var schema = new ParquetSchema(
new DataField<string>("encounter_id"),
new DataField<string>("patient_id"),
new DataField<string>("previous_status"),
new DataField<string>("new_status"),
new DataField<string>("changed_at"),
new DataField<int>("kafka_partition"),
new DataField<long>("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();
}
}