feature: Elasticsearch CQRS Projection and Analytics Endpoints
This commit is contained in:
@@ -0,0 +1,94 @@
|
||||
using Elastic.Clients.Elasticsearch;
|
||||
using Elastic.Clients.Elasticsearch.IndexManagement;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
public class ElasticIndexProvisioner : IHostedService
|
||||
{
|
||||
private readonly ElasticsearchClient _elastic;
|
||||
private readonly ElasticsearchOptions _options;
|
||||
private readonly ILogger<ElasticIndexProvisioner> _logger;
|
||||
|
||||
public ElasticIndexProvisioner(
|
||||
ElasticsearchClient elastic,
|
||||
IOptions<ElasticsearchOptions> options,
|
||||
ILogger<ElasticIndexProvisioner> logger)
|
||||
{
|
||||
_elastic = elastic;
|
||||
_options = options.Value;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
public async Task StartAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
await EnsureIndexAsync<PatientEncounterDocument>(
|
||||
_options.Indices.PatientEncounters, BuildPatientEncountersMapping());
|
||||
await EnsureIndexAsync<ObservationDocument>(
|
||||
_options.Indices.Observations, BuildObservationsMapping());
|
||||
await EnsureIndexAsync<ClinicalAlertDocument>(
|
||||
_options.Indices.ClinicalAlerts, BuildClinicalAlertsMapping());
|
||||
}
|
||||
|
||||
private async Task EnsureIndexAsync<T>(
|
||||
string indexName,
|
||||
Action<CreateIndexRequestDescriptor<T>> configure) where T : class
|
||||
{
|
||||
var exists = await _elastic.Indices.ExistsAsync(indexName);
|
||||
if (exists.Exists)
|
||||
{
|
||||
_logger.LogInformation("Elasticsearch index '{Index}' already exists — skipping", indexName);
|
||||
return;
|
||||
}
|
||||
|
||||
var resp = await _elastic.Indices.CreateAsync<T>(indexName, configure);
|
||||
if (!resp.IsValidResponse)
|
||||
throw new InvalidOperationException(
|
||||
$"Failed to create Elasticsearch index '{indexName}': {resp.DebugInformation}");
|
||||
|
||||
_logger.LogInformation("Created Elasticsearch index '{Index}'", indexName);
|
||||
}
|
||||
|
||||
private Action<CreateIndexRequestDescriptor<PatientEncounterDocument>> BuildPatientEncountersMapping() =>
|
||||
d => d.Mappings(m => m.Properties(p => p
|
||||
.Keyword(k => k.EncounterId)
|
||||
.Keyword(k => k.PatientId)
|
||||
.Keyword(k => k.Mrn)
|
||||
// text for full-text search + keyword sub-field for exact sort/filter
|
||||
.Text(t => t.PatientName, tf => tf
|
||||
.Fields(f => f.Keyword(k => k.PatientName)))
|
||||
.Keyword(k => k.Department)
|
||||
.Keyword(k => k.Status)
|
||||
.Keyword(k => k.AttendingPhysician)
|
||||
.Date(d => d.AdmittedAt)
|
||||
.IntegerNumber(i => i.OpenAlertCount)
|
||||
.Date(d => d.LastObservationAt!)
|
||||
));
|
||||
|
||||
private Action<CreateIndexRequestDescriptor<ObservationDocument>> BuildObservationsMapping() =>
|
||||
d => d.Mappings(m => m.Properties(p => p
|
||||
.Keyword(k => k.ObservationId)
|
||||
.Keyword(k => k.EncounterId)
|
||||
.Keyword(k => k.PatientId)
|
||||
.Keyword(k => k.Mrn)
|
||||
.Keyword(k => k.ObservationCode)
|
||||
// float: observations are decimal values like 97.3, 2.5, 118.0
|
||||
// auto-mapped 'long' would truncate fractional parts silently
|
||||
.FloatNumber(f => f.Value)
|
||||
.Keyword(k => k.Unit)
|
||||
.Keyword(k => k.Source)
|
||||
.Date(d => d.RecordedAt)
|
||||
));
|
||||
|
||||
private Action<CreateIndexRequestDescriptor<ClinicalAlertDocument>> BuildClinicalAlertsMapping() =>
|
||||
d => d.Mappings(m => m.Properties(p => p
|
||||
.Keyword(k => k.AlertId)
|
||||
.Keyword(k => k.EncounterId)
|
||||
.Keyword(k => k.PatientId)
|
||||
.Keyword(k => k.Department)
|
||||
.Keyword(k => k.AlertType)
|
||||
.Keyword(k => k.Severity)
|
||||
.Keyword(k => k.Status)
|
||||
.Date(d => d.TriggeredAt)
|
||||
));
|
||||
|
||||
public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask;
|
||||
}
|
||||
@@ -0,0 +1,241 @@
|
||||
using System.Text.Json;
|
||||
using Confluent.Kafka;
|
||||
using Elastic.Clients.Elasticsearch;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
public class EsIndexerService : BackgroundService
|
||||
{
|
||||
private static readonly JsonSerializerOptions EventJsonOptions = new()
|
||||
{
|
||||
PropertyNameCaseInsensitive = true
|
||||
};
|
||||
|
||||
private readonly ElasticsearchClient _elastic;
|
||||
private readonly KafkaOptions _kafkaOptions;
|
||||
private readonly ElasticsearchOptions _esOptions;
|
||||
private readonly ILogger<EsIndexerService> _logger;
|
||||
|
||||
public EsIndexerService(
|
||||
ElasticsearchClient elastic,
|
||||
IOptions<KafkaOptions> kafkaOptions,
|
||||
IOptions<ElasticsearchOptions> esOptions,
|
||||
ILogger<EsIndexerService> logger)
|
||||
{
|
||||
_elastic = elastic;
|
||||
_kafkaOptions = kafkaOptions.Value;
|
||||
_esOptions = esOptions.Value;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
||||
{
|
||||
var config = new ConsumerConfig
|
||||
{
|
||||
BootstrapServers = _kafkaOptions.BootstrapServers,
|
||||
GroupId = "es-indexer",
|
||||
AutoOffsetReset = AutoOffsetReset.Earliest,
|
||||
// Manual commit: offset is only committed after a successful Elasticsearch write.
|
||||
// If the process crashes between ES write and commit, the message is reprocessed.
|
||||
// Consumers must be idempotent. See idempotency contract above.
|
||||
EnableAutoCommit = false,
|
||||
EnablePartitionEof = false
|
||||
};
|
||||
|
||||
using var consumer = new ConsumerBuilder<string, string>(config).Build();
|
||||
|
||||
consumer.Subscribe(new[]
|
||||
{
|
||||
_kafkaOptions.Topics.ObservationRecorded,
|
||||
_kafkaOptions.Topics.AlertGenerated,
|
||||
_kafkaOptions.Topics.EncounterStatusChanged
|
||||
});
|
||||
|
||||
_logger.LogInformation("EsIndexerService started. Subscribed to 3 topics.");
|
||||
|
||||
try
|
||||
{
|
||||
while (!stoppingToken.IsCancellationRequested)
|
||||
{
|
||||
ConsumeResult<string, string>? result = null;
|
||||
try
|
||||
{
|
||||
result = consumer.Consume(stoppingToken);
|
||||
await DispatchAsync(result.Topic, result.Message.Value, stoppingToken);
|
||||
consumer.Commit(result);
|
||||
}
|
||||
catch (OperationCanceledException)
|
||||
{
|
||||
break;
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogError(ex,
|
||||
"EsIndexer failed processing topic={Topic} offset={Offset} — not committing",
|
||||
result?.Topic, result?.Offset.Value);
|
||||
// Do not commit: message will be redelivered on restart
|
||||
await Task.Delay(1000, stoppingToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
finally
|
||||
{
|
||||
consumer.Close();
|
||||
}
|
||||
}
|
||||
|
||||
private Task DispatchAsync(string topic, string payload, CancellationToken ct) => topic switch
|
||||
{
|
||||
var t when t == _kafkaOptions.Topics.EncounterStatusChanged =>
|
||||
HandleEncounterStatusChangedAsync(payload, ct),
|
||||
var t when t == _kafkaOptions.Topics.ObservationRecorded =>
|
||||
HandleObservationRecordedAsync(payload, ct),
|
||||
var t when t == _kafkaOptions.Topics.AlertGenerated =>
|
||||
HandleAlertGeneratedAsync(payload, ct),
|
||||
_ => Task.CompletedTask
|
||||
};
|
||||
|
||||
// --- encounter.status.changed ---
|
||||
// Upserts the patient_encounters document. DocAsUpsert=true means:
|
||||
// if the document does not exist, it is created; if it exists, it is replaced.
|
||||
// Idempotent: applying the same event twice produces the same document.
|
||||
private async Task HandleEncounterStatusChangedAsync(string payload, CancellationToken ct)
|
||||
{
|
||||
var evt = JsonSerializer.Deserialize<EncounterStatusChangedEvent>(payload, EventJsonOptions)!;
|
||||
|
||||
var doc = new PatientEncounterDocument
|
||||
{
|
||||
EncounterId = evt.EncounterId.ToString(),
|
||||
PatientId = evt.PatientId.ToString(),
|
||||
Mrn = evt.Mrn,
|
||||
PatientName = evt.PatientName,
|
||||
Department = evt.Department,
|
||||
Status = evt.NewStatus,
|
||||
AttendingPhysician = evt.AttendingPhysician,
|
||||
AdmittedAt = evt.AdmittedAt,
|
||||
OpenAlertCount = 0,
|
||||
LastObservationAt = null
|
||||
};
|
||||
|
||||
var resp = await _elastic.UpdateAsync<PatientEncounterDocument, PatientEncounterDocument>(
|
||||
_esOptions.Indices.PatientEncounters,
|
||||
evt.EncounterId.ToString(),
|
||||
u => u.Doc(doc).DocAsUpsert(true),
|
||||
ct);
|
||||
|
||||
if (!resp.IsValidResponse)
|
||||
throw new InvalidOperationException(
|
||||
$"ES upsert failed for encounter {evt.EncounterId}: {resp.DebugInformation}");
|
||||
|
||||
_logger.LogDebug("Upserted patient_encounters for encounter {Id} → status={Status}",
|
||||
evt.EncounterId, evt.NewStatus);
|
||||
}
|
||||
|
||||
// --- observation.recorded ---
|
||||
// Indexes the observation by observationId — idempotent PUT.
|
||||
// Also updates lastObservationAt on the parent encounter document using a script.
|
||||
private async Task HandleObservationRecordedAsync(string payload, CancellationToken ct)
|
||||
{
|
||||
var evt = JsonSerializer.Deserialize<ObservationRecordedEvent>(payload, EventJsonOptions)!;
|
||||
|
||||
// Index into observations — document ID is observationId
|
||||
var doc = new ObservationDocument
|
||||
{
|
||||
ObservationId = evt.ObservationId.ToString(),
|
||||
EncounterId = evt.EncounterId.ToString(),
|
||||
PatientId = evt.PatientId.ToString(),
|
||||
Mrn = evt.Mrn ?? string.Empty,
|
||||
ObservationCode = evt.ObservationCode,
|
||||
Value = (double)evt.Value,
|
||||
Unit = evt.Unit,
|
||||
Source = evt.Source,
|
||||
RecordedAt = evt.RecordedAt
|
||||
};
|
||||
|
||||
var indexResp = await _elastic.IndexAsync(
|
||||
doc,
|
||||
i => i.Index(_esOptions.Indices.Observations).Id(doc.ObservationId),
|
||||
ct);
|
||||
|
||||
if (!indexResp.IsValidResponse)
|
||||
throw new InvalidOperationException(
|
||||
$"ES index failed for observation {evt.ObservationId}: {indexResp.DebugInformation}");
|
||||
|
||||
// Update lastObservationAt on patient_encounters using a conditional script:
|
||||
// only update if the new recordedAt is later than the stored value.
|
||||
// This handles out-of-order delivery: an older observation re-processed after a
|
||||
// newer one must not overwrite lastObservationAt with a stale timestamp.
|
||||
var updateResp = await _elastic.UpdateAsync<PatientEncounterDocument, object>(
|
||||
_esOptions.Indices.PatientEncounters,
|
||||
evt.EncounterId.ToString(),
|
||||
u => u
|
||||
.Script(new Script(new InlineScript
|
||||
{
|
||||
Source = """
|
||||
if (ctx._source.lastObservationAt == null ||
|
||||
params.recordedAt > ctx._source.lastObservationAt) {
|
||||
ctx._source.lastObservationAt = params.recordedAt;
|
||||
}
|
||||
""",
|
||||
Language = ScriptLanguage.Painless,
|
||||
Params = new Dictionary<string, object>
|
||||
{
|
||||
["recordedAt"] = evt.RecordedAt.ToString("O")
|
||||
}
|
||||
}))
|
||||
.RetryOnConflict(3),
|
||||
ct);
|
||||
|
||||
if (!updateResp.IsValidResponse && updateResp.Result != Result.NotFound)
|
||||
_logger.LogWarning(
|
||||
"Could not update lastObservationAt for encounter {Id} — encounter document may not exist yet",
|
||||
evt.EncounterId);
|
||||
}
|
||||
|
||||
// --- alert.generated ---
|
||||
// Indexes the alert by alertId. Also increments openAlertCount on patient_encounters.
|
||||
// TRADE-OFF: openAlertCount increment is not idempotent for partial reprocessing.
|
||||
// It is correct for full replay from offset 0 (the stated recovery procedure).
|
||||
// For production, use a set of counted alert IDs in the script to enforce idempotency.
|
||||
private async Task HandleAlertGeneratedAsync(string payload, CancellationToken ct)
|
||||
{
|
||||
var evt = JsonSerializer.Deserialize<AlertGeneratedEvent>(payload, EventJsonOptions)!;
|
||||
|
||||
var doc = new ClinicalAlertDocument
|
||||
{
|
||||
AlertId = evt.AlertId.ToString(),
|
||||
EncounterId = evt.EncounterId.ToString(),
|
||||
PatientId = evt.PatientId.ToString(),
|
||||
Department = evt.Department ?? string.Empty,
|
||||
AlertType = evt.AlertType,
|
||||
Severity = evt.Severity,
|
||||
Status = "Open",
|
||||
TriggeredAt = evt.TriggeredAt
|
||||
};
|
||||
|
||||
var indexResp = await _elastic.IndexAsync(
|
||||
doc,
|
||||
i => i.Index(_esOptions.Indices.ClinicalAlerts).Id(doc.AlertId),
|
||||
ct);
|
||||
|
||||
if (!indexResp.IsValidResponse)
|
||||
throw new InvalidOperationException(
|
||||
$"ES index failed for alert {evt.AlertId}: {indexResp.DebugInformation}");
|
||||
|
||||
// Increment openAlertCount on the parent encounter document
|
||||
var updateResp = await _elastic.UpdateAsync<PatientEncounterDocument, object>(
|
||||
_esOptions.Indices.PatientEncounters,
|
||||
evt.EncounterId.ToString(),
|
||||
u => u
|
||||
.Script(new Script(new InlineScript
|
||||
{
|
||||
Source = "ctx._source.openAlertCount += 1",
|
||||
Language = ScriptLanguage.Painless
|
||||
}))
|
||||
.RetryOnConflict(3),
|
||||
ct);
|
||||
|
||||
if (!updateResp.IsValidResponse && updateResp.Result != Result.NotFound)
|
||||
_logger.LogWarning(
|
||||
"Could not increment openAlertCount for encounter {Id}", evt.EncounterId);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user