1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
| public class KafkaService : IKafkaService
{
public async Task SubscribeAsync<TMessage>(IEnumerable<string> topics, Action<TMessage> handler,
CancellationToken cancellationToken) where TMessage : class
{
var config = new ConsumerConfig
{
BootstrapServers = "127.0.0.1:9092",
GroupId = "app-consumers",
EnableAutoCommit = false,
AutoOffsetReset = AutoOffsetReset.Earliest
};
using var consumer = new ConsumerBuilder<Ignore, string>(config)
.SetErrorHandler((_, e) => Console.WriteLine($"Error: {e.Reason}"))
.SetPartitionsAssignedHandler((c, partitions) =>
Console.WriteLine($"Assigned: {string.Join(", ", partitions)}"))
.Build();
consumer.Subscribe(topics);
try
{
while (!cancellationToken.IsCancellationRequested)
{
var result = consumer.Consume(cancellationToken);
var message = JsonConvert.DeserializeObject<TMessage>(result.Message.Value);
handler(message);
consumer.Commit(result);
}
}
catch (OperationCanceledException)
{
consumer.Close();
}
}
}
|