Skip to content

Commit 20c1377

Browse files
authored
[Feature]Add support to measure event/stream performance metrics (#166)
* [feature]add support to measure event/stream performance metrics * remove the parameter when init construct of StreamPressureOptions
1 parent 40659fb commit 20c1377

18 files changed

Lines changed: 1352 additions & 9 deletions

Directory.Packages.props

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@
1717
<PackageVersion Include="Azure.AI.TextAnalytics" Version="5.3.0" />
1818
<PackageVersion Include="JetBrains.Annotations" Version="2024.2.0" />
1919
<PackageVersion Include="Microsoft.EntityFrameworkCore.Sqlite.Core" Version="8.0.4" />
20-
<PackageVersion Include="Microsoft.Extensions.Options" Version="8.0.2" />
20+
<PackageVersion Include="Microsoft.Extensions.Options" Version="9.0.0" />
2121
<PackageVersion Include="Microsoft.Orleans.Serialization.Abstractions" Version="9.0.1" />
2222
<PackageVersion Include="Microsoft.SemanticKernel" Version="1.33.0" />
2323
<PackageVersion Include="Microsoft.SemanticKernel.Core" Version="1.33.0" />
@@ -99,7 +99,8 @@
9999
<PackageVersion Include="Microsoft.CodeCoverage" Version="17.9.0" />
100100
<PackageVersion Include="Microsoft.Extensions.AI" Version="9.0.1-preview.1.24570.5" />
101101
<PackageVersion Include="Microsoft.Extensions.AI.Abstractions" Version="9.0.1-preview.1.24570.5" />
102-
<PackageVersion Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="8.0.2" />
102+
<PackageVersion Include="Microsoft.Extensions.DependencyInjection" Version="9.0.0" />
103+
<PackageVersion Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="9.0.0" />
103104
<PackageVersion Include="Microsoft.Extensions.FileProviders.Embedded" Version="8.0.4" />
104105
<PackageVersion Include="Microsoft.Extensions.Http.Polly" Version="9.0.0" />
105106
<PackageVersion Include="Microsoft.Extensions.Hosting" Version="9.0.0" />

src/Aevatar.Core.Abstractions/Events/EventWrapperBase.cs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,9 @@ public abstract class EventWrapperBase
1818
[Id(0)]
1919
public Dictionary<string, string> ContextMetadata { get; set; } = new Dictionary<string, string>();
2020

21+
[Id(1)]
22+
public DateTime PublishedTimestampUtc { get; set; }
23+
2124
protected EventWrapperBase()
2225
{
2326
// Initialize ContextMetadata

src/Aevatar.Core/Aevatar.Core.csproj

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,12 @@
1717

1818
<ItemGroup>
1919
<PackageReference Include="Microsoft.Orleans.EventSourcing" />
20+
<PackageReference Include="Microsoft.Orleans.Streaming" />
21+
<PackageReference Include="Orleans.Streams.Kafka" />
22+
<PackageReference Include="Microsoft.Extensions.DependencyInjection" />
23+
<PackageReference Include="Microsoft.Extensions.DependencyInjection.Abstractions" />
24+
<PackageReference Include="Microsoft.Extensions.Logging" />
25+
<PackageReference Include="Microsoft.Extensions.Options" />
2026
<PackageReference Include="AElf.OpenTelemetry" />
2127
</ItemGroup>
2228

src/Aevatar.Core/AevatarSyncWorker.cs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,9 @@ protected sealed override async Task<TResponse> PerformWork(TRequest request,
2929
{
3030
var response = await PerformLongRunTask(request);
3131
_logger.LogInformation($"Performed long run task for request of type {typeof(TRequest).FullName}, response is {response}");
32-
await _asyncStream.OnNextAsync(new EventWrapper<TResponse>(response, Guid.NewGuid(), this.GetGrainId()));
32+
var eventWrapper = new EventWrapper<TResponse>(response, Guid.NewGuid(), this.GetGrainId());
33+
eventWrapper.PublishedTimestampUtc = DateTime.UtcNow;
34+
await _asyncStream.OnNextAsync(eventWrapper);
3335
return response;
3436
}
3537
catch (Exception ex)

src/Aevatar.Core/BroadCastGAgentBase.cs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,7 @@ public async Task BroadCastEventAsync<T>(string streamIdString, T @event) where
6868
{
6969
var stream = GenStream<T>(streamIdString);
7070
var eventWrapper = new EventWrapper<T>(@event, Guid.NewGuid(), this.GetGrainId());
71+
eventWrapper.PublishedTimestampUtc = DateTime.UtcNow;
7172
await stream.OnNextAsync(eventWrapper);
7273
}
7374

src/Aevatar.Core/EventWrapperBaseAsyncObserver.cs

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,17 @@ public static EventWrapperBaseAsyncObserver Create(Func<EventWrapperBase, Task>
3838

3939
public async Task OnNextAsync(EventWrapperBase item, StreamSequenceToken? token = null)
4040
{
41+
if (item == null)
42+
{
43+
throw new ArgumentNullException(nameof(item));
44+
}
45+
46+
var eventType = item.GetType().Name;
47+
var latency = (DateTime.UtcNow - item.PublishedTimestampUtc).TotalSeconds;
48+
_logger?.LogInformation("[OnNextAsync] Consuming event_type={EventType} PublishedTimestampUtc={PublishedTimestampUtc} latency={Latency}s",
49+
eventType, item.PublishedTimestampUtc, latency);
50+
Observability.EventPublishLatencyMetrics.Record(latency, item, MethodName, ParameterTypeName, null);
51+
4152
await _func(item);
4253
}
4354

src/Aevatar.Core/GAgentAsyncObserver.cs

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,19 +2,29 @@
22
using Orleans.Streams;
33
using OrleansCodeGen.Orleans.Runtime;
44
using System.Diagnostics;
5+
using System.Diagnostics.Metrics;
6+
using Microsoft.Extensions.Logging;
57

68
namespace Aevatar.Core;
79

810
public class GAgentAsyncObserver : IAsyncObserver<EventWrapperBase>
911
{
1012
private readonly List<EventWrapperBaseAsyncObserver> _observers;
1113
private readonly string _grainId;
14+
private readonly ILogger<GAgentAsyncObserver> _logger;
15+
private static readonly Meter Meter = new(OpenTelemetryConstants.AevatarStreamsMeterName);
16+
private static readonly Histogram<double> PublishLatencyHistogram = Meter.CreateHistogram<double>(
17+
OpenTelemetryConstants.EventPublishLatencyHistogram, "s", "Event publish-to-consume latency");
1218

13-
public GAgentAsyncObserver(List<EventWrapperBaseAsyncObserver> observers, string grainId)
19+
public GAgentAsyncObserver(List<EventWrapperBaseAsyncObserver> observers, string grainId, ILogger<GAgentAsyncObserver> logger)
1420
{
1521
_observers = observers;
1622
_grainId = grainId;
23+
_logger = logger;
1724
}
25+
26+
public GAgentAsyncObserver(List<EventWrapperBaseAsyncObserver> observers, string grainId)
27+
: this(observers, grainId, null) { }
1828

1929
/// <summary>
2030
/// Finds observers that match the given event type
@@ -65,8 +75,9 @@ private async Task BroadcastToObservers(Func<EventWrapperBaseAsyncObserver, Task
6575

6676
public async Task OnNextAsync(EventWrapperBase item, StreamSequenceToken? token = null)
6777
{
68-
// Extract event and ID from wrapper
6978
var (eventType, eventId) = EventWrapperHelper.ExtractProperties(item);
79+
var latency = (DateTime.UtcNow - item.PublishedTimestampUtc).TotalSeconds;
80+
Observability.EventPublishLatencyMetrics.Record(latency, item, null, null, _logger);
7081

7182
// Try to create an activity with parent context if available
7283
var activity = ActivityHelper.CreateEventActivity(item, eventType, _grainId, eventId, token);

src/Aevatar.Core/GAgentBase.Observers.cs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -103,8 +103,11 @@ private EventWrapperBaseAsyncObserver CreateMethodObserver(
103103
?.GetValue(item)!;
104104
var grainId = (GrainId)item.GetType().GetProperty(nameof(EventWrapper<TEvent>.GrainId))
105105
?.GetValue(item)!;
106+
var publishedTimestamp = (DateTime)item.GetType().GetProperty(nameof(EventWrapper<TEvent>.PublishedTimestampUtc))
107+
?.GetValue(item)!;
106108

107109
var eventWrapper = new EventWrapper<TEvent>(eventType, eventId, grainId);
110+
eventWrapper.PublishedTimestampUtc = publishedTimestamp;
108111

109112
Logger.LogInformation("Handling event {EventWrapper} in method {MethodName}", eventWrapper,
110113
method.Name);
@@ -260,6 +263,7 @@ private async Task PublishResponse(EventBase result, Guid eventId)
260263
(TEvent)result,
261264
eventId,
262265
this.GetGrainId());
266+
responseWrapper.PublishedTimestampUtc = DateTime.UtcNow;
263267

264268
await PublishAsync(responseWrapper);
265269
}

src/Aevatar.Core/GAgentBase.Publish.cs

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@ protected async Task<Guid> PublishAsync<T>(T @event) where T : EventBase
3737
{
3838
// Create event wrapper with context propagation
3939
var eventWrapper = new EventWrapper<T>(@event, eventId, this.GetGrainId());
40+
eventWrapper.PublishedTimestampUtc = DateTime.UtcNow;
4041
if (State.Parent == null)
4142
{
4243
Logger.LogInformation(
@@ -66,7 +67,9 @@ private async Task PublishEventUpwardsAsync<T>(T @event, Guid eventId) where T :
6667
{
6768
try
6869
{
69-
await SendEventUpwardsAsync(new EventWrapper<T>(@event, eventId, GrainId));
70+
var eventWrapper = new EventWrapper<T>(@event, eventId, GrainId);
71+
eventWrapper.PublishedTimestampUtc = DateTime.UtcNow;
72+
await SendEventUpwardsAsync(eventWrapper);
7073
Logger.LogDebug("{GrainId} published {Event} to upwards", GrainId.ToString(),
7174
JsonConvert.SerializeObject(@event));
7275
}
@@ -94,7 +97,7 @@ private async Task SendEventUpwardsAsync<T>(EventWrapper<T> eventWrapper) where
9497
{
9598
Logger.LogError("{GrainId} failed to send event {EventWrapper} upwards", GrainId.ToString(),
9699
eventWrapper);
97-
throw new EventPublishingException($"{GrainId.ToString()} failed to event event upwards", ex);
100+
throw new EventPublishingException($"{GrainId.ToString()} failed to send event upwards", ex);
98101
}
99102
}
100103

@@ -111,7 +114,7 @@ private async Task SendEventToSelfAsync<T>(EventWrapper<T> eventWrapper) where T
111114
{
112115
Logger.LogError("{GrainId} failed to send event {EventWrapper} to itself", GrainId.ToString(),
113116
eventWrapper);
114-
throw new EventPublishingException($"{GrainId.ToString()} failed to event event to itself", ex);
117+
throw new EventPublishingException($"{GrainId.ToString()} failed to send event to itself", ex);
115118
}
116119
}
117120

@@ -135,7 +138,7 @@ private async Task SendEventDownwardsAsync<T>(EventWrapper<T> eventWrapper) wher
135138
{
136139
Logger.LogError("{GrainId} failed to send event {EventWrapper} downwards", GrainId.ToString(),
137140
eventWrapper);
138-
throw new EventPublishingException($"{GrainId.ToString()} failed to event event downwards", ex);
141+
throw new EventPublishingException($"{GrainId.ToString()} failed to send event downwards", ex);
139142
}
140143
}
141144
}
Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
1+
using System.Diagnostics.Metrics;
2+
using Aevatar.Core.Abstractions;
3+
using Microsoft.Extensions.Logging;
4+
5+
namespace Aevatar.Core.Observability;
6+
7+
public static class EventPublishLatencyMetrics
8+
{
9+
private static readonly Meter Meter = new(OpenTelemetryConstants.AevatarStreamsMeterName);
10+
private static readonly Histogram<double> PublishLatencyHistogram = Meter.CreateHistogram<double>(
11+
OpenTelemetryConstants.EventPublishLatencyHistogram, "s", "Event publish-to-consume latency");
12+
13+
public static void Record(double latency, EventWrapperBase item, string? methodName, string? parameterTypeName, ILogger? logger = null)
14+
{
15+
var eventType = item.GetType().Name;
16+
var grainId = item.GetType().GetProperty("GrainId")?.GetValue(item)?.ToString() ?? "unknown";
17+
var eventId = item.GetType().GetProperty("EventId")?.GetValue(item)?.ToString() ?? "unknown";
18+
PublishLatencyHistogram.Record(latency,
19+
new KeyValuePair<string, object?>("grain_id", grainId),
20+
new KeyValuePair<string, object?>("method_name", methodName ?? "unknown"),
21+
new KeyValuePair<string, object?>("parameter_type", parameterTypeName ?? "unknown"));
22+
logger?.LogInformation("[PublishLatency] latency={Latency}s event_type={EventType} grain_id={GrainId} event_id={EventId} method={MethodName} parameter={ParameterType}",
23+
latency, eventType, grainId, eventId, methodName ?? "unknown", parameterTypeName ?? "unknown");
24+
}
25+
}

0 commit comments

Comments
 (0)