test phase 23 verification script
This commit is contained in:
@@ -59,7 +59,11 @@ public class TrendAnalyzerService : BackgroundService
|
||||
evt.RecordedAt,
|
||||
stoppingToken);
|
||||
|
||||
if (outcome.Outcome == TrendOutcome.RapidDeterioration)
|
||||
if (outcome.Outcome == TrendOutcome.EncounterNotFound)
|
||||
_logger.LogWarning(
|
||||
"Skipping stale observation.recorded event — encounter={Id} code={Code} offset={Offset}",
|
||||
evt.EncounterId, outcome.ObservationCode, result.Offset.Value);
|
||||
else if (outcome.Outcome == TrendOutcome.RapidDeterioration)
|
||||
_logger.LogInformation(
|
||||
"RAPID_DETERIORATION alert via consumer — encounter={Id} code={Code} rate={Rate}/min",
|
||||
evt.EncounterId, outcome.ObservationCode, outcome.RatePerMinute);
|
||||
|
||||
@@ -0,0 +1,137 @@
|
||||
using System.Text.Json;
|
||||
|
||||
/// <summary>
|
||||
/// Best-effort deserializers for Kafka payloads written to Parquet.
|
||||
/// Tolerates legacy/minimal observation shapes (e.g. test outbox rows using
|
||||
/// <c>code</c> instead of <c>observationCode</c>).
|
||||
/// </summary>
|
||||
public static class DataLakeEventParser
|
||||
{
|
||||
public static ObservationRow ParseObservationRow(string payload, long offset, int partition)
|
||||
{
|
||||
var d = JsonDocument.Parse(payload).RootElement;
|
||||
var encounterId = GetString(d, "encounterId");
|
||||
var observationId = GetString(d, "observationId");
|
||||
if (string.IsNullOrEmpty(observationId))
|
||||
observationId = $"legacy:{encounterId}:{offset}";
|
||||
|
||||
return new ObservationRow(
|
||||
ObservationId : observationId,
|
||||
EncounterId : encounterId,
|
||||
PatientId : GetString(d, "patientId"),
|
||||
Mrn : GetString(d, "mrn"),
|
||||
ObservationCode : GetString(d, "observationCode", "code"),
|
||||
Value : GetDouble(d, "value"),
|
||||
Unit : GetString(d, "unit"),
|
||||
Source : GetString(d, "source"),
|
||||
RecordedAt : GetTimestampString(d, "recordedAt"),
|
||||
KafkaPartition : partition,
|
||||
KafkaOffset : offset);
|
||||
}
|
||||
|
||||
public static AlertRow ParseAlertRow(string payload, long offset, int partition)
|
||||
{
|
||||
var d = JsonDocument.Parse(payload).RootElement;
|
||||
return new AlertRow(
|
||||
AlertId : GetString(d, "alertId"),
|
||||
EncounterId : GetString(d, "encounterId"),
|
||||
PatientId : GetString(d, "patientId"),
|
||||
AlertType : GetString(d, "alertType"),
|
||||
Severity : GetString(d, "severity"),
|
||||
Details : GetString(d, "details"),
|
||||
TriggeredAt : GetTimestampString(d, "triggeredAt"),
|
||||
KafkaPartition : partition,
|
||||
KafkaOffset : offset);
|
||||
}
|
||||
|
||||
public static EncounterStatusRow ParseEncounterRow(string payload, long offset, int partition)
|
||||
{
|
||||
var d = JsonDocument.Parse(payload).RootElement;
|
||||
return new EncounterStatusRow(
|
||||
EncounterId : GetString(d, "encounterId"),
|
||||
PatientId : GetString(d, "patientId"),
|
||||
PreviousStatus : GetString(d, "previousStatus"),
|
||||
NewStatus : GetString(d, "newStatus"),
|
||||
ChangedAt : GetTimestampString(d, "changedAt"),
|
||||
KafkaPartition : partition,
|
||||
KafkaOffset : offset);
|
||||
}
|
||||
|
||||
public static string ExtractDatePath(string topic, string payload, KafkaTopicOptions topics)
|
||||
{
|
||||
try
|
||||
{
|
||||
var doc = JsonDocument.Parse(payload);
|
||||
var ts = topic switch
|
||||
{
|
||||
var t when t == topics.ObservationRecorded => GetTimestamp(doc.RootElement, "recordedAt"),
|
||||
var t when t == topics.AlertGenerated => GetTimestamp(doc.RootElement, "triggeredAt"),
|
||||
var t when t == topics.EncounterStatusChanged => GetTimestamp(doc.RootElement, "changedAt"),
|
||||
_ => DateTimeOffset.UtcNow,
|
||||
};
|
||||
return $"{ts.Year:D4}/{ts.Month:D2}/{ts.Day:D2}";
|
||||
}
|
||||
catch
|
||||
{
|
||||
var now = DateTimeOffset.UtcNow;
|
||||
return $"{now.Year:D4}/{now.Month:D2}/{now.Day:D2}";
|
||||
}
|
||||
}
|
||||
|
||||
private static string GetString(JsonElement d, string primary, string? alternate = null)
|
||||
{
|
||||
if (TryGetProperty(d, primary, out var prop))
|
||||
return ElementToString(prop);
|
||||
if (alternate is not null && TryGetProperty(d, alternate, out prop))
|
||||
return ElementToString(prop);
|
||||
return "";
|
||||
}
|
||||
|
||||
private static double GetDouble(JsonElement d, string name)
|
||||
{
|
||||
if (!TryGetProperty(d, name, out var prop))
|
||||
return 0;
|
||||
|
||||
return prop.ValueKind switch
|
||||
{
|
||||
JsonValueKind.Number => prop.GetDouble(),
|
||||
JsonValueKind.String => double.TryParse(prop.GetString(), out var v) ? v : 0,
|
||||
_ => 0
|
||||
};
|
||||
}
|
||||
|
||||
private static string GetTimestampString(JsonElement d, string name)
|
||||
{
|
||||
if (!TryGetProperty(d, name, out var prop))
|
||||
return "";
|
||||
|
||||
if (prop.ValueKind == JsonValueKind.String)
|
||||
return prop.GetString() ?? "";
|
||||
|
||||
return prop.TryGetDateTimeOffset(out var dt) ? dt.ToString("O") : "";
|
||||
}
|
||||
|
||||
private static DateTimeOffset GetTimestamp(JsonElement d, string name)
|
||||
{
|
||||
if (!TryGetProperty(d, name, out var prop))
|
||||
return DateTimeOffset.UtcNow;
|
||||
|
||||
if (prop.ValueKind == JsonValueKind.String)
|
||||
return prop.GetDateTimeOffset();
|
||||
|
||||
return prop.TryGetDateTimeOffset(out var dt) ? dt : DateTimeOffset.UtcNow;
|
||||
}
|
||||
|
||||
private static bool TryGetProperty(JsonElement d, string name, out JsonElement prop) =>
|
||||
d.TryGetProperty(name, out prop) && prop.ValueKind != JsonValueKind.Null;
|
||||
|
||||
private static string ElementToString(JsonElement prop) =>
|
||||
prop.ValueKind switch
|
||||
{
|
||||
JsonValueKind.String => prop.GetString() ?? "",
|
||||
JsonValueKind.Number => prop.GetRawText(),
|
||||
_ when prop.ValueKind == JsonValueKind.True || prop.ValueKind == JsonValueKind.False
|
||||
=> prop.GetBoolean().ToString(),
|
||||
_ => prop.GetRawText()
|
||||
};
|
||||
}
|
||||
@@ -191,22 +191,28 @@ public sealed class DataLakeWriterService : BackgroundService
|
||||
}
|
||||
}
|
||||
|
||||
private string ExtractDatePath(string topic, string payload) =>
|
||||
DataLakeEventParser.ExtractDatePath(topic, payload, _kafkaOptions.Topics);
|
||||
|
||||
private async Task<byte[]> BuildParquetAsync(
|
||||
string topic, List<BufferedEvent> events, int partition)
|
||||
{
|
||||
if (topic == _kafkaOptions.Topics.ObservationRecorded)
|
||||
{
|
||||
var rows = events.Select(e => ParseObservationRow(e, partition)).ToList();
|
||||
var rows = events.Select(e => DataLakeEventParser.ParseObservationRow(
|
||||
e.Payload, e.Offset, partition)).ToList();
|
||||
return await ParquetFileBuilder.BuildObservationsAsync(rows);
|
||||
}
|
||||
if (topic == _kafkaOptions.Topics.AlertGenerated)
|
||||
{
|
||||
var rows = events.Select(e => ParseAlertRow(e, partition)).ToList();
|
||||
var rows = events.Select(e => DataLakeEventParser.ParseAlertRow(
|
||||
e.Payload, e.Offset, partition)).ToList();
|
||||
return await ParquetFileBuilder.BuildAlertsAsync(rows);
|
||||
}
|
||||
if (topic == _kafkaOptions.Topics.EncounterStatusChanged)
|
||||
{
|
||||
var rows = events.Select(e => ParseEncounterRow(e, partition)).ToList();
|
||||
var rows = events.Select(e => DataLakeEventParser.ParseEncounterRow(
|
||||
e.Payload, e.Offset, partition)).ToList();
|
||||
return await ParquetFileBuilder.BuildEncountersAsync(rows);
|
||||
}
|
||||
throw new InvalidOperationException($"Unknown topic: {topic}");
|
||||
@@ -233,80 +239,6 @@ public sealed class DataLakeWriterService : BackgroundService
|
||||
return $"{folder}/{key.DatePath}/partition-{key.Partition}-offset-{firstOffset:D10}.parquet";
|
||||
}
|
||||
|
||||
private string ExtractDatePath(string topic, string payload)
|
||||
{
|
||||
try
|
||||
{
|
||||
var doc = JsonDocument.Parse(payload);
|
||||
var ts = topic switch
|
||||
{
|
||||
var t when t == _kafkaOptions.Topics.ObservationRecorded => doc.RootElement.GetProperty("recordedAt").GetDateTimeOffset(),
|
||||
var t when t == _kafkaOptions.Topics.AlertGenerated => doc.RootElement.GetProperty("triggeredAt").GetDateTimeOffset(),
|
||||
var t when t == _kafkaOptions.Topics.EncounterStatusChanged => doc.RootElement.GetProperty("changedAt").GetDateTimeOffset(),
|
||||
_ => DateTimeOffset.UtcNow,
|
||||
};
|
||||
return $"{ts.Year:D4}/{ts.Month:D2}/{ts.Day:D2}";
|
||||
}
|
||||
catch
|
||||
{
|
||||
// Malformed payload: use today so the event is not lost.
|
||||
var now = DateTimeOffset.UtcNow;
|
||||
return $"{now.Year:D4}/{now.Month:D2}/{now.Day:D2}";
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Payload deserializers — each reads only the fields needed for the Parquet row.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
private static ObservationRow ParseObservationRow(BufferedEvent e, int partition)
|
||||
{
|
||||
var d = JsonDocument.Parse(e.Payload).RootElement;
|
||||
return new ObservationRow(
|
||||
ObservationId : d.GetProperty("observationId").GetString() ?? "",
|
||||
EncounterId : d.GetProperty("encounterId").GetString() ?? "",
|
||||
PatientId : d.GetProperty("patientId").GetString() ?? "",
|
||||
Mrn : d.TryGetProperty("mrn", out var mrn) ? mrn.GetString() ?? "" : "",
|
||||
ObservationCode : d.GetProperty("observationCode").GetString() ?? "",
|
||||
Value : d.GetProperty("value").GetDouble(),
|
||||
Unit : d.GetProperty("unit").GetString() ?? "",
|
||||
Source : d.GetProperty("source").GetString() ?? "",
|
||||
RecordedAt : d.GetProperty("recordedAt").GetString() ?? "",
|
||||
KafkaPartition : partition,
|
||||
KafkaOffset : e.Offset
|
||||
);
|
||||
}
|
||||
|
||||
private static AlertRow ParseAlertRow(BufferedEvent e, int partition)
|
||||
{
|
||||
var d = JsonDocument.Parse(e.Payload).RootElement;
|
||||
return new AlertRow(
|
||||
AlertId : d.GetProperty("alertId").GetString() ?? "",
|
||||
EncounterId : d.GetProperty("encounterId").GetString() ?? "",
|
||||
PatientId : d.GetProperty("patientId").GetString() ?? "",
|
||||
AlertType : d.GetProperty("alertType").GetString() ?? "",
|
||||
Severity : d.GetProperty("severity").GetString() ?? "",
|
||||
Details : d.TryGetProperty("details", out var det) ? det.GetString() ?? "" : "",
|
||||
TriggeredAt : d.GetProperty("triggeredAt").GetString() ?? "",
|
||||
KafkaPartition : partition,
|
||||
KafkaOffset : e.Offset
|
||||
);
|
||||
}
|
||||
|
||||
private static EncounterStatusRow ParseEncounterRow(BufferedEvent e, int partition)
|
||||
{
|
||||
var d = JsonDocument.Parse(e.Payload).RootElement;
|
||||
return new EncounterStatusRow(
|
||||
EncounterId : d.GetProperty("encounterId").GetString() ?? "",
|
||||
PatientId : d.GetProperty("patientId").GetString() ?? "",
|
||||
PreviousStatus : d.GetProperty("previousStatus").GetString() ?? "",
|
||||
NewStatus : d.GetProperty("newStatus").GetString() ?? "",
|
||||
ChangedAt : d.GetProperty("changedAt").GetString() ?? "",
|
||||
KafkaPartition : partition,
|
||||
KafkaOffset : e.Offset
|
||||
);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// MinIO upload
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
@@ -4,5 +4,6 @@ public enum TrendOutcome
|
||||
InsufficientHistory,
|
||||
Stable,
|
||||
RapidDeterioration,
|
||||
AlertAlreadyOpen
|
||||
AlertAlreadyOpen,
|
||||
EncounterNotFound
|
||||
}
|
||||
@@ -76,6 +76,17 @@ public class TrendDetector
|
||||
if (!TrendCalculator.ExceedsThreshold(observationCode, rate.Value, threshold))
|
||||
return new TrendResult(TrendOutcome.Stable, observationCode, rate);
|
||||
|
||||
using (var scope = _services.CreateScope())
|
||||
{
|
||||
var db = scope.ServiceProvider.GetRequiredService<AppDbContext>();
|
||||
if (!await db.Encounters.AnyAsync(e => e.Id == encounterId, ct))
|
||||
{
|
||||
_logger.LogWarning(
|
||||
"Skipping trend alert for unknown encounter {EncounterId}", encounterId);
|
||||
return new TrendResult(TrendOutcome.EncounterNotFound, observationCode, rate);
|
||||
}
|
||||
}
|
||||
|
||||
var details = TrendCalculator.DescribeTrend(observationCode, rate.Value, value);
|
||||
var created = await TryCreateAlertAsync(
|
||||
encounterId, patientId, observationCode, details, rate.Value, ct);
|
||||
|
||||
Reference in New Issue
Block a user