Browse Source

Performance improved.

pull/65/head
Sebastian Stehle 9 years ago
parent
commit
822fb30813
  1. 53
      src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs
  2. 2
      tests/Benchmarks/Program.cs
  3. 5
      tests/Benchmarks/Tests/HandleEvents.cs

53
src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs

@ -15,10 +15,9 @@ using MongoDB.Bson;
using MongoDB.Driver; using MongoDB.Driver;
using Squidex.Infrastructure.CQRS.Events; using Squidex.Infrastructure.CQRS.Events;
// ReSharper disable ConvertIfStatementToConditionalTernaryExpression
// ReSharper disable ClassNeverInstantiated.Local
// ReSharper disable UnusedMember.Local
// ReSharper disable InvertIf // ReSharper disable InvertIf
// ReSharper disable ConvertIfStatementToConditionalTernaryExpression
// ReSharper disable TooWideLocalVariableScope
namespace Squidex.Infrastructure.MongoDb.EventStore namespace Squidex.Infrastructure.MongoDb.EventStore
{ {
@ -48,7 +47,7 @@ namespace Squidex.Infrastructure.MongoDb.EventStore
protected override async Task SetupCollectionAsync(IMongoCollection<MongoEventCommit> collection) protected override async Task SetupCollectionAsync(IMongoCollection<MongoEventCommit> collection)
{ {
await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.EventStream).Ascending(x => x.Timestamp)); await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.Timestamp).Ascending(x => x.EventStream));
await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.EventStreamOffset).Ascending(x => x.EventStream), new CreateIndexOptions { Unique = true }); await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.EventStreamOffset).Ascending(x => x.EventStream), new CreateIndexOptions { Unique = true });
} }
@ -70,20 +69,36 @@ namespace Squidex.Infrastructure.MongoDb.EventStore
Guard.NotNull(callback, nameof(callback)); Guard.NotNull(callback, nameof(callback));
var tokenTimestamp = EmptyTimestamp; var tokenTimestamp = EmptyTimestamp;
var tokenEventStreamNumber = -1; var tokenCommitSize = -1;
var tokenCommitOffset = -1;
var isEndOfCommit = false;
if (position != null) if (position != null)
{ {
var token = ParsePosition(position); var token = ParsePosition(position);
tokenTimestamp = token.Timestamp; tokenTimestamp = token.Timestamp;
tokenEventStreamNumber = token.EventStreamNumber; tokenCommitSize = token.CommitSize;
tokenCommitOffset = token.CommitOffset;
isEndOfCommit = tokenCommitOffset == tokenCommitSize - 1;
if (isEndOfCommit)
{
tokenCommitOffset = -1;
}
}
var filters = new List<FilterDefinition<MongoEventCommit>>();
if (isEndOfCommit)
{
filters.Add(Filter.Gt(x => x.Timestamp, tokenTimestamp));
} }
else
var filters = new List<FilterDefinition<MongoEventCommit>>
{ {
Filter.Gte(x => x.Timestamp, tokenTimestamp) filters.Add(Filter.Gte(x => x.Timestamp, tokenTimestamp));
}; }
if (!string.IsNullOrWhiteSpace(streamFilter) && !string.Equals(streamFilter, "*", StringComparison.OrdinalIgnoreCase)) if (!string.IsNullOrWhiteSpace(streamFilter) && !string.Equals(streamFilter, "*", StringComparison.OrdinalIgnoreCase))
{ {
@ -110,20 +125,22 @@ namespace Squidex.Infrastructure.MongoDb.EventStore
await Collection.Find(filter).SortBy(x => x.Timestamp).ForEachAsync(async commit => await Collection.Find(filter).SortBy(x => x.Timestamp).ForEachAsync(async commit =>
{ {
var isGreaterTimestamp = commit.Timestamp > tokenTimestamp;
var eventStreamNumber = (int)commit.EventStreamOffset; var eventStreamNumber = (int)commit.EventStreamOffset;
var commitOffset = 0;
foreach (var e in commit.Events) foreach (var e in commit.Events)
{ {
eventStreamNumber++; eventStreamNumber++;
if (isGreaterTimestamp || eventStreamNumber > tokenEventStreamNumber) if (commitOffset > tokenCommitOffset)
{ {
var eventData = new EventData { EventId = e.EventId, Metadata = e.Metadata, Payload = e.Payload, Type = e.Type }; var eventData = new EventData { EventId = e.EventId, Metadata = e.Metadata, Payload = e.Payload, Type = e.Type };
var eventToken = CreateToken(commit.Timestamp, eventStreamNumber); var eventToken = CreateToken(commit.Timestamp, commitOffset, commit.Events.Length);
await callback(new StoredEvent(eventToken, eventStreamNumber, eventData)); await callback(new StoredEvent(eventToken, eventStreamNumber, eventData));
commitOffset++;
} }
else else
{ {
@ -201,18 +218,18 @@ namespace Squidex.Infrastructure.MongoDb.EventStore
return -1; return -1;
} }
private static string CreateToken(BsonTimestamp timestamp, int eventStreamNumber) private static string CreateToken(BsonTimestamp timestamp, int commitOffset, int commitSize)
{ {
var parts = new object[] { timestamp.Timestamp, timestamp.Increment, eventStreamNumber }; var parts = new object[] { timestamp.Timestamp, timestamp.Increment, commitOffset, commitSize };
return string.Join("-", parts); return string.Join("-", parts);
} }
private static (BsonTimestamp Timestamp, int EventStreamNumber) ParsePosition(string position) private static (BsonTimestamp Timestamp, int CommitOffset, int CommitSize) ParsePosition(string position)
{ {
var parts = position.Split('-'); var parts = position.Split('-');
return (new BsonTimestamp(int.Parse(parts[0]), int.Parse(parts[1])), int.Parse(parts[2])); return (new BsonTimestamp(int.Parse(parts[0]), int.Parse(parts[1])), int.Parse(parts[2]), int.Parse(parts[3]));
} }
} }
} }

2
tests/Benchmarks/Program.cs

@ -25,7 +25,7 @@ namespace Benchmarks
public static void Main(string[] args) public static void Main(string[] args)
{ {
var id = args.Length > 0 ? args[0] : "handleEvents"; var id = args.Length > 0 ? args[0] : string.Empty;
var benchmark = Benchmarks.Find(x => x.Id == id); var benchmark = Benchmarks.Find(x => x.Id == id);

5
tests/Benchmarks/Tests/HandleEvents.cs

@ -9,6 +9,7 @@
using System; using System;
using System.Collections.Generic; using System.Collections.Generic;
using System.Threading.Tasks; using System.Threading.Tasks;
using MongoDB.Bson;
using MongoDB.Driver; using MongoDB.Driver;
using Newtonsoft.Json; using Newtonsoft.Json;
using Squidex.Infrastructure; using Squidex.Infrastructure;
@ -82,7 +83,7 @@ namespace Benchmarks.Tests
private readonly TypeNameRegistry typeNameRegistry = new TypeNameRegistry().Map(typeof(MyEvent)); private readonly TypeNameRegistry typeNameRegistry = new TypeNameRegistry().Map(typeof(MyEvent));
private readonly EventDataFormatter formatter; private readonly EventDataFormatter formatter;
private readonly JsonSerializerSettings serializerSettings = new JsonSerializerSettings(); private readonly JsonSerializerSettings serializerSettings = new JsonSerializerSettings();
private const int NumEvents = 1000; private const int NumEvents = 5000;
private IMongoClient mongoClient; private IMongoClient mongoClient;
private IMongoDatabase mongoDatabase; private IMongoDatabase mongoDatabase;
private IEventStore eventStore; private IEventStore eventStore;
@ -143,7 +144,7 @@ namespace Benchmarks.Tests
if (eventConsumer.EventNumbers.Count != NumEvents) if (eventConsumer.EventNumbers.Count != NumEvents)
{ {
throw new InvalidOperationException($"Only {eventConsumer.EventNumbers.Count} have been handled"); throw new InvalidOperationException($"{eventConsumer.EventNumbers.Count} Events have been handled");
} }
for (var i = 0; i < eventConsumer.EventNumbers.Count; i++) for (var i = 0; i < eventConsumer.EventNumbers.Count; i++)

Loading…
Cancel
Save