fix: Shutdown causes issues in datalake and minio

This commit is contained in:
voltsrage
2026-06-18 16:28:58 +08:00
parent ddde7fee31
commit c3fbc20ddc
3 changed files with 23 additions and 5 deletions
@@ -89,9 +89,11 @@ public sealed class DataLakeWriterService : BackgroundService
finally finally
{ {
// Final flush on shutdown so buffered events are not lost. // Final flush on shutdown so buffered events are not lost.
// Use CancellationToken.None — host shutdown cancels the execute token before
// MinIO uploads finish, which surfaces as TaskCanceledException noise.
if (_buffer.Values.Sum(v => v.Count) > 0) if (_buffer.Values.Sum(v => v.Count) > 0)
{ {
try { await FlushAsync(consumer, ct); } try { await FlushAsync(consumer, CancellationToken.None); }
catch (Exception ex) catch (Exception ex)
{ {
_logger.LogError(ex, "DataLakeWriter shutdown flush failed — some events may be re-read on next start"); _logger.LogError(ex, "DataLakeWriter shutdown flush failed — some events may be re-read on next start");
@@ -143,19 +145,33 @@ public sealed class DataLakeWriterService : BackgroundService
{ {
// Log and continue — a failed file for one key must not prevent other // Log and continue — a failed file for one key must not prevent other
// keys from flushing. The uncommitted offsets will cause reprocessing. // keys from flushing. The uncommitted offsets will cause reprocessing.
if (ex is OperationCanceledException && ct.IsCancellationRequested)
{
_logger.LogInformation(
"[DATA-LAKE] Flush canceled for key {Key} during shutdown", key);
}
else
{
_logger.LogError(ex, "[DATA-LAKE] Failed to write file for key {Key}", key); _logger.LogError(ex, "[DATA-LAKE] Failed to write file for key {Key}", key);
} }
} }
}
// Commit only after all files are uploaded. // Commit only after at least one file uploaded successfully.
// Events for any key that failed above will be re-read on next startup. // Events for any key that failed above will be re-read on next startup.
if (_highWatermarks.Any()) if (filesWritten > 0 && _highWatermarks.Any())
{ {
consumer.Commit(_highWatermarks.Values); consumer.Commit(_highWatermarks.Values);
_logger.LogInformation( _logger.LogInformation(
"[DATA-LAKE] Committed offsets for {PartitionCount} partitions after flushing {FileCount} files", "[DATA-LAKE] Committed offsets for {PartitionCount} partitions after flushing {FileCount} files",
_highWatermarks.Count, filesWritten); _highWatermarks.Count, filesWritten);
} }
else if (filesWritten == 0 && _highWatermarks.Any())
{
_logger.LogWarning(
"[DATA-LAKE] Skipping offset commit — no files were written ({BufferedPartitions} partitions buffered)",
_highWatermarks.Count);
}
_buffer.Clear(); _buffer.Clear();
_highWatermarks.Clear(); _highWatermarks.Clear();
@@ -256,7 +272,7 @@ public sealed class DataLakeWriterService : BackgroundService
PatientId : d.GetProperty("patientId").GetString() ?? "", PatientId : d.GetProperty("patientId").GetString() ?? "",
AlertType : d.GetProperty("alertType").GetString() ?? "", AlertType : d.GetProperty("alertType").GetString() ?? "",
Severity : d.GetProperty("severity").GetString() ?? "", Severity : d.GetProperty("severity").GetString() ?? "",
Details : d.GetProperty("details").GetString() ?? "", Details : d.TryGetProperty("details", out var det) ? det.GetString() ?? "" : "",
TriggeredAt : d.GetProperty("triggeredAt").GetString() ?? "", TriggeredAt : d.GetProperty("triggeredAt").GetString() ?? "",
KafkaPartition : partition, KafkaPartition : partition,
KafkaOffset : e.Offset KafkaOffset : e.Offset
@@ -145,6 +145,7 @@ public class SirsDetector
patientId, patientId,
alertType = AlertType.SepsisWarning.ToDbString(), alertType = AlertType.SepsisWarning.ToDbString(),
severity = "Critical", severity = "Critical",
details,
triggeredAt, triggeredAt,
partitionKey = encounterId.ToString() partitionKey = encounterId.ToString()
}), }),
@@ -111,6 +111,7 @@ public class WarningEvaluator
patientId, patientId,
alertType = alertType.ToDbString(), alertType = alertType.ToDbString(),
severity = "Warning", severity = "Warning",
details,
triggeredAt, triggeredAt,
partitionKey = encounterId.ToString() partitionKey = encounterId.ToString()
}), }),