Add machine and retry to Redis Event message

This commit is contained in:
Haiping Chen 2024-12-15 16:27:37 +00:00
parent 78ae361d3f
commit b3bdc01fc8
2 changed files with 36 additions and 18 deletions

View file

@ -39,10 +39,7 @@ public class RedisPublisher : IEventPublisher
// Add a message to the stream, keeping only the latest 1 million messages
var messageId = await db.StreamAddAsync(channel,
[
new NameValueEntry("message", message),
new NameValueEntry("timestamp", DateTime.UtcNow.ToString("o"))
],
AssembleMessage(message),
maxLength: 1000 * 10000);
_logger.LogInformation($"Published message {channel} {message} ({messageId})");
@ -82,6 +79,17 @@ public class RedisPublisher : IEventPublisher
return exists;
}
private NameValueEntry[] AssembleMessage(RedisValue message, int retry = 0)
{
return
[
new NameValueEntry("message", message),
new NameValueEntry("timestamp", DateTime.UtcNow.ToString("o")),
new NameValueEntry("machine", Environment.MachineName),
new NameValueEntry("retry", retry),
];
}
public async Task ReDispatchAsync(string channel, int count = 10, string order = "asc")
{
var db = _redis.GetDatabase();
@ -93,10 +101,12 @@ public class RedisPublisher : IEventPublisher
try
{
var messageId = await db.StreamAddAsync(channel, [
new NameValueEntry("message", entry.Values[0].Value),
new NameValueEntry("timestamp", DateTime.UtcNow.ToString("o"))
]);
var message = entry.Values.First(x => x.Name == "message").Value;
var retryKv = entry.Values.FirstOrDefault(x => x.Name == "retry");
int.TryParse(retryKv.Value, out int retry);
var messageId = await db.StreamAddAsync(channel,
AssembleMessage(message, retry: retry + 1),
maxLength: 1000 * 10000);
_logger.LogWarning($"ReDispatched message: {channel} {entry.Values[0].Value} ({messageId})");

View file

@ -58,23 +58,31 @@ public class RedisSubscriber : IEventSubscriber
break;
}
if (priorityEnabled)
try
{
if (await HandleGroupMessage(db, $"{channel}-{EventPriority.High}", group, consumer, received) > 0)
if (priorityEnabled)
{
continue;
}
if (await HandleGroupMessage(db, $"{channel}-{EventPriority.High}", group, consumer, received) > 0)
{
continue;
}
if (await HandleGroupMessage(db, $"{channel}-{EventPriority.Medium}", group, consumer, received) > 0)
if (await HandleGroupMessage(db, $"{channel}-{EventPriority.Medium}", group, consumer, received) > 0)
{
continue;
}
await HandleGroupMessage(db, $"{channel}-{EventPriority.Low}", group, consumer, received);
}
else
{
continue;
await HandleGroupMessage(db, channel, group, consumer, received);
}
await HandleGroupMessage(db, $"{channel}-{EventPriority.Low}", group, consumer, received);
}
else
catch (Exception ex)
{
await HandleGroupMessage(db, channel, group, consumer, received);
_logger.LogError($"Error processing message: {ex.Message}\r\n{ex}");
await Task.Delay(1000 * 60);
}
}
}