-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathProgram.cs
More file actions
57 lines (52 loc) · 2.03 KB
/
Copy pathProgram.cs
File metadata and controls
57 lines (52 loc) · 2.03 KB
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
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
using Confluent.Kafka;
using KafkaSimpleProducerSubscriber.KafkaOperators;
using Microsoft.Extensions.Configuration;
using System.Threading.Tasks;
namespace KafkaSimpleProducerSubscriber
{
public class Program
{
/// <summary>
/// Gets the configuration.
/// </summary>
/// <value>The configuration.</value>
public IConfiguration Configuration { get; }
public Program(IConfiguration configuration)
{
Configuration = configuration;
}
static async Task Main()
{
var topic = "OrderEventQA2";
var bootstrapServers = "127.0.0.1:9092";
CreateKafkaTopics.Create(false);
PublishMessageToKafka publishMessageToKafka = KafkaProducerBuilder(bootstrapServers);
await publishMessageToKafka.Publish(topic);
ConsumeMessageFromKafka consumeMessageFromKafka = KafkaConsumerBuilder(bootstrapServers);
consumeMessageFromKafka.Consume(topic);
}
private static ConsumeMessageFromKafka KafkaConsumerBuilder(string bootstrapServers)
{
var config = new ConsumerConfig
{
GroupId = "ConsumerGroup01",
BootstrapServers = bootstrapServers,
AutoOffsetReset = AutoOffsetReset.Earliest
};
var kafkaConsumer = new ConsumerBuilder<int, string>(config).Build();
var consumeMessageFromKafka = new ConsumeMessageFromKafka(kafkaConsumer);
return consumeMessageFromKafka;
}
private static PublishMessageToKafka KafkaProducerBuilder(string bootstrapServers)
{
var producerConfig = new ProducerConfig
{
BootstrapServers = bootstrapServers,
Acks = Acks.All
};
var kafkaProducer = new ProducerBuilder<int, string>(producerConfig).Build();
var publishMessageToKafka = new PublishMessageToKafka(kafkaProducer);
return publishMessageToKafka;
}
}
}