using System.Text.Json; using Confluent.Kafka; using Microsoft.Extensions.Options; using Minio.DataModel.Args; public sealed class DataLakeWriterService : BackgroundService { private readonly KafkaOptions _kafkaOptions; private readonly DataLakeOptions _opts; private readonly MinioOptions _minioOpts; private readonly ILogger _logger; // Buffer key: identifies one Parquet file-to-be. // Events sharing a topic, date, and Kafka partition land in the same file. private record BufferKey(string Topic, string DatePath, int Partition); private record BufferedEvent(string Payload, long Offset); private readonly Dictionary> _buffer = new(); // Track the highest offset per topic-partition for post-flush commit. private readonly Dictionary _highWatermarks = new(); public DataLakeWriterService( IOptions kafkaOptions, IOptions opts, IOptions minioOpts, ILogger logger) { _kafkaOptions = kafkaOptions.Value; _opts = opts.Value; _minioOpts = minioOpts.Value; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken ct) { var consumerConfig = new ConsumerConfig { BootstrapServers = _kafkaOptions.BootstrapServers, GroupId = "data-lake-writer", AutoOffsetReset = AutoOffsetReset.Earliest, EnableAutoCommit = false, }; using var consumer = new ConsumerBuilder(consumerConfig).Build(); consumer.Subscribe(new[] { _kafkaOptions.Topics.ObservationRecorded, _kafkaOptions.Topics.AlertGenerated, _kafkaOptions.Topics.EncounterStatusChanged, }); _logger.LogInformation( "DataLakeWriterService started. FlushCount={FlushCount} FlushIntervalSeconds={FlushInterval}", _opts.FlushCount, _opts.FlushIntervalSeconds); var lastFlush = DateTimeOffset.UtcNow; try { while (!ct.IsCancellationRequested) { ConsumeResult? result; try { result = consumer.Consume(TimeSpan.FromMilliseconds(500)); } catch (ConsumeException ex) { _logger.LogError(ex, "DataLakeWriter consume error"); continue; } if (result is not null) AddToBuffer(result); var totalBuffered = _buffer.Values.Sum(v => v.Count); var shouldFlushCount = totalBuffered >= _opts.FlushCount; var shouldFlushTime = DateTimeOffset.UtcNow - lastFlush >= TimeSpan.FromSeconds(_opts.FlushIntervalSeconds); if ((shouldFlushCount || shouldFlushTime) && totalBuffered > 0) { await FlushAsync(consumer, ct); lastFlush = DateTimeOffset.UtcNow; } } } finally { // Final flush on shutdown so buffered events are not lost. if (_buffer.Values.Sum(v => v.Count) > 0) { try { await FlushAsync(consumer, ct); } catch (Exception ex) { _logger.LogError(ex, "DataLakeWriter shutdown flush failed — some events may be re-read on next start"); } } consumer.Close(); } } private void AddToBuffer(ConsumeResult result) { var datePath = ExtractDatePath(result.Topic, result.Message.Value); var key = new BufferKey(result.Topic, datePath, result.Partition.Value); if (!_buffer.TryGetValue(key, out var list)) { list = new List(); _buffer[key] = list; } list.Add(new BufferedEvent(result.Message.Value, result.Offset.Value)); // Track highest offset per topic-partition for post-flush commit. var tp = new TopicPartition(result.Topic, result.Partition); _highWatermarks[tp] = new TopicPartitionOffset(tp, result.Offset + 1); } private async Task FlushAsync(IConsumer consumer, CancellationToken ct) { var filesWritten = 0; foreach (var (key, events) in _buffer) { if (events.Count == 0) continue; try { var firstOffset = events.Min(e => e.Offset); var objectKey = BuildObjectKey(key, firstOffset); var bytes = await BuildParquetAsync(key.Topic, events, key.Partition); await UploadToMinioAsync(objectKey, bytes, ct); filesWritten++; _logger.LogInformation( "[DATA-LAKE] Wrote {Count} events → {ObjectKey} ({Bytes} bytes)", events.Count, objectKey, bytes.Length); } catch (Exception ex) { // Log and continue — a failed file for one key must not prevent other // keys from flushing. The uncommitted offsets will cause reprocessing. _logger.LogError(ex, "[DATA-LAKE] Failed to write file for key {Key}", key); } } // Commit only after all files are uploaded. // Events for any key that failed above will be re-read on next startup. if (_highWatermarks.Any()) { consumer.Commit(_highWatermarks.Values); _logger.LogInformation( "[DATA-LAKE] Committed offsets for {PartitionCount} partitions after flushing {FileCount} files", _highWatermarks.Count, filesWritten); } _buffer.Clear(); _highWatermarks.Clear(); } private async Task BuildParquetAsync( string topic, List events, int partition) { if (topic == _kafkaOptions.Topics.ObservationRecorded) { var rows = events.Select(e => ParseObservationRow(e, partition)).ToList(); return await ParquetFileBuilder.BuildObservationsAsync(rows); } if (topic == _kafkaOptions.Topics.AlertGenerated) { var rows = events.Select(e => ParseAlertRow(e, partition)).ToList(); return await ParquetFileBuilder.BuildAlertsAsync(rows); } if (topic == _kafkaOptions.Topics.EncounterStatusChanged) { var rows = events.Select(e => ParseEncounterRow(e, partition)).ToList(); return await ParquetFileBuilder.BuildEncountersAsync(rows); } throw new InvalidOperationException($"Unknown topic: {topic}"); } // --------------------------------------------------------------------------- // Object key and date partition helpers // --------------------------------------------------------------------------- // File path: observations/2025/01/15/partition-0-offset-0000001000.parquet // The date comes from the event timestamp, not the wall clock. // Events from the same encounter that cross midnight are written into the date // bucket matching their recorded_at timestamp — consistent with how Athena and // Spark partition-prune by event time, not ingest time. private string BuildObjectKey(BufferKey key, long firstOffset) { var folder = key.Topic switch { var t when t == _kafkaOptions.Topics.ObservationRecorded => "observations", var t when t == _kafkaOptions.Topics.AlertGenerated => "alerts", var t when t == _kafkaOptions.Topics.EncounterStatusChanged => "encounters", _ => "unknown", }; 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.GetProperty("details").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 // --------------------------------------------------------------------------- private async Task UploadToMinioAsync(string objectKey, byte[] bytes, CancellationToken ct) { var client = MinioClientFactory.Build(_minioOpts); var bucket = _opts.BucketName; var exists = await client.BucketExistsAsync( new BucketExistsArgs().WithBucket(bucket), ct); if (!exists) await client.MakeBucketAsync(new MakeBucketArgs().WithBucket(bucket), ct); await client.PutObjectAsync(new PutObjectArgs() .WithBucket(bucket) .WithObject(objectKey) .WithStreamData(new MemoryStream(bytes)) .WithObjectSize(bytes.Length) .WithContentType("application/octet-stream"), ct); } }