.net core C#集成kafka,生产者和消费者
·
准备依赖包:Confluent.Kafka
生产者
1、在appsettings.json中配置kafka地址等配置信息
"Kafka": {
"BootstrapServers": "localhost:9092",
"ClientId": "my-app"
}
2、定义接口和实现类
接口:
namespace kafkaweb.Service
{
public interface IKafkaProducerService
{
string produceMessage(string topic,string key, string message);
}
}
实现类:
using Confluent.Kafka;
using static Confluent.Kafka.ConfigPropertyNames;
namespace kafkaweb.Service.Impl
{
public class KafkaProducerServiceImpl : IKafkaProducerService
{
private readonly IProducer<string, string> producer;
private readonly IConfiguration _configuration;
public KafkaProducerServiceImpl(IConfiguration configuration)
{
_configuration = configuration;
ProducerConfig config = new ProducerConfig { BootstrapServers = _configuration["Kafka:BootstrapServers"], ClientId= _configuration["Kafka:ClientId"] };
producer = new ProducerBuilder<string, string>(config).Build();
}
public string produceMessage(string topic,string key, string message)
{
producer.ProduceAsync(topic, new Message<string, string> { Key = key, Value = message });
return "OK";
}
}
}
3、在Program.cs中进行注入
//注入服务
builder.Services.AddSingleton<IKafkaProducerService, KafkaProducerServiceImpl>();
4、在控制器中测试使用
[ApiController]
[Route("[controller]")]
public class KafkaController : ControllerBase
{
private readonly IKafkaProducerService kafkaProducerService;
public KafkaController(IKafkaProducerService _kafkaProducerService)
{
kafkaProducerService = _kafkaProducerService;
}
[HttpPost(Name = "/sendMessage")]
public string sendMessage([FromForm] Dictionary<string,string> dict)
{
kafkaProducerService.produceMessage(dict["topic"], dict["key"], dict["message"]);
return "OK";
}
}
消费者
1、在appsettings.json中配置kafka地址等配置信息
"Kafka": {
"BootstrapServers": "localhost:9092",
"ClientId": "zmtest",
"GroupId": "my-consumer"
}
2、定义接口和实现类
接口
public interface IKafkaConsumerService
{
string consumerKafka(string topic);
}
实现类
public class KafkaConsumerServiceImpl : IKafkaConsumerService
{
private readonly IConsumer<string, string> consumer;
private readonly IConfiguration _configuration;
public KafkaConsumerServiceImpl(IConfiguration configuration)
{
_configuration = configuration;
ConsumerConfig consumerConfig = new ConsumerConfig() { BootstrapServers = _configuration["Kafka:BootstrapServers"], GroupId = _configuration["Kafka:GroupId"], EnableAutoCommit = false, AutoOffsetReset = AutoOffsetReset.Earliest };
consumer=new ConsumerBuilder<string, string>(consumerConfig).Build();
}
public void executeConsume(string topic)
{
//订阅
consumer.Subscribe(topic);
while (true)
{
ConsumeResult<string, string> consumeResult = consumer.Consume(1000);
Console.WriteLine("测试数据");
if (consumeResult != null) {
Console.WriteLine(consumeResult.Message.Key);
Console.WriteLine(consumeResult.Message.Value);
//提交位移
consumer.Commit(consumeResult);
}
}
}
public string consumerKafka(string topic)
{
Thread thread=new Thread(() =>
{
executeConsume(topic);
});
thread.IsBackground = true;
thread.Start();
return "OK";
}
}
3、在Program.cs中进行注入
builder.Services.AddSingleton<IKafkaConsumerService,KafkaConsumerServiceImpl>();
4、在控制器中测试使用
//消费
[HttpGet(Name = "consumeMessage")]
public string consumeMessage()
{
return kafkaConsumerService.consumerKafka("test");
}
更多推荐




所有评论(0)