准备依赖包: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");
}
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐