.net core kafka 入门实例 一篇看懂
时间:2020-05-20
本文章向大家介绍.net core kafka 入门实例 一篇看懂,主要包括.net core kafka 入门实例 一篇看懂使用实例、应用技巧、基本知识点总结和需要注意事项,具有一定的参考价值,需要的朋友可以参考一下。
/// <summary> /// 消费者 /// </summary> public interface IKafkaConsumer : IDisposable { /// <summary> /// 消费数据 /// </summary> /// <typeparam name="T"></typeparam> /// <returns></returns> T Consume<T>() where T : class; } public interface IKafkaProducer : IDisposable { /// <summary> /// 发布消息 /// </summary> /// <typeparam name="T"></typeparam> /// <param name="key"></param> /// <param name="data"></param> /// <param name="operateType"></param> /// <returns></returns> bool Produce<T>(string key, T data, int operateType) where T : class; }
实现方法
using Confluent.Kafka; using Newtonsoft.Json; using System; using System.Collections.Generic; using System.Text; namespace Kafka { public class KafkaConsumer : IKafkaConsumer { private bool disposeHasBeenCalled = false; private readonly object disposeHasBeenCalledLockObj = new object(); private readonly IConsumer<string, string> _consumer; /// <summary> /// 构造函数,初始化配置 /// </summary> /// <param name="config">配置参数</param> /// <param name="topic">主题名称</param> public KafkaConsumer(ConsumerConfig config, string topic) { _consumer = new ConsumerBuilder<string, string>(config).Build(); _consumer.Subscribe(topic); } /// <summary> /// 消费 /// </summary> /// <returns></returns> public T Consume<T>() where T : class { try { var result = _consumer.Consume(TimeSpan.FromSeconds(1)); if (result != null) { if (typeof(T) == typeof(string)) return (T)Convert.ChangeType(result.Value, typeof(T)); return JsonConvert.DeserializeObject<T>(result.Value); } } catch (ConsumeException e) { Console.WriteLine($"consume error: {e.Error.Reason}"); } catch (Exception e) { Console.WriteLine($"consume error: {e.Message}"); } return default; } /// <summary> /// 释放 /// </summary> public void Dispose() { Dispose(true); GC.SuppressFinalize(this); } /// <summary> /// Dispose /// </summary> /// <param name="disposing"></param> protected virtual void Dispose(bool disposing) { lock (disposeHasBeenCalledLockObj) { if (disposeHasBeenCalled) { return; } disposeHasBeenCalled = true; } if (disposing) { _consumer?.Close(); } } }
}
public class KafkaProducer : IKafkaProducer { private bool disposeHasBeenCalled = false; private readonly object disposeHasBeenCalledLockObj = new object(); private readonly IProducer<string, string> _producer; private readonly string _topic; /// <summary> /// 构造函数,初始化配置 /// </summary> /// <param name="config">配置参数</param> /// <param name="topic">主题名称</param> public KafkaProducer(ProducerConfig config, string topic) { _producer = new ProducerBuilder<string, string>(config).Build(); _topic = topic; } /// <summary> /// 发布消息 /// </summary> /// <typeparam name="T">数据实体</typeparam> /// <param name="key">数据key,partition分区会根据key</param> /// <param name="data">数据</param> /// <param name="operateType">操作类型[增、删、改等不同类型]</param> /// <returns></returns> public bool Produce<T>(string key, T data, int operateType) where T : class { var obj = JsonConvert.SerializeObject(new { Type = operateType, Data = data }); try { var result = _producer.ProduceAsync(_topic, new Message<string, string> { Key = key, Value = obj }).ConfigureAwait(false).GetAwaiter().GetResult(); #if DEBUG Console.WriteLine($"Topic: {result.Topic} Partition: {result.Partition} Offset: {result.Offset}"); #endif return true; } catch (ProduceException<string, string> e) { Console.WriteLine($"Delivery failed: {e.Error.Reason}"); } catch (Exception e) { Console.WriteLine($"Delivery failed: {e.Message}"); } return false; } /// <summary> /// 释放 /// </summary> public void Dispose() { Dispose(true); GC.SuppressFinalize(this); } /// <summary> /// Dispose /// </summary> /// <param name="disposing"></param> protected virtual void Dispose(bool disposing) { lock (disposeHasBeenCalledLockObj) { if (disposeHasBeenCalled) { return; } disposeHasBeenCalled = true; } if (disposing) { _producer?.Dispose(); } } }
static void Main(string[] args) { var config = new ProducerConfig { BootstrapServers = "localhost:9092", Acks = Acks.All }; //发送消息 using (var kafkaProducer = new KafkaProducer(config, "topic-d")) { var result = kafkaProducer.Produce<object>("a", new { name = "猪八戒3" }, 1); } Console.WriteLine("消息发送成功"); } static void Main(string[] args) { var config = new ConsumerConfig { BootstrapServers = "localhost:9092", GroupId = "test", AutoOffsetReset = AutoOffsetReset.Earliest }; string text; Console.WriteLine("接受中......"); while ((text = Console.ReadLine()) != "q") { //接受消息 using (var kafkaProducer = new KafkaConsumer(config, "topic-d")) { var result = kafkaProducer.Consume<object>(); if (result != null) { Console.WriteLine(result.ToString()); } } } }
上结果
可以看到,消息已经收到了。这个demo里,消费端要一直处于正常状态才行,才能消费生产者得信息
如果您觉得阅读本文对您有帮助,请点一下“推荐”按钮,您的“推荐”将是我最大的写作动力!
本文版权归作者和博客园共有,来源网址:https://www.cnblogs.com/DanielYao/欢迎各位转载,但是未经作者本人同意,转载文章之后必须在文章页面明显位置给出作者和原文连接,否则保留追究法律责任的权利。
原文地址:https://www.cnblogs.com/DanielYao/p/12922745.html
- JavaScript 教程
- JavaScript 编辑工具
- JavaScript 与HTML
- JavaScript 与Java
- JavaScript 数据结构
- JavaScript 基本数据类型
- JavaScript 特殊数据类型
- JavaScript 运算符
- JavaScript typeof 运算符
- JavaScript 表达式
- JavaScript 类型转换
- JavaScript 基本语法
- JavaScript 注释
- Javascript 基本处理流程
- Javascript 选择结构
- Javascript if 语句
- Javascript if 语句的嵌套
- Javascript switch 语句
- Javascript 循环结构
- Javascript 循环结构实例
- Javascript 跳转语句
- Javascript 控制语句总结
- Javascript 函数介绍
- Javascript 函数的定义
- Javascript 函数调用
- Javascript 几种特殊的函数
- JavaScript 内置函数简介
- Javascript eval() 函数
- Javascript isFinite() 函数
- Javascript isNaN() 函数
- parseInt() 与 parseFloat()
- escape() 与 unescape()
- Javascript 字符串介绍
- Javascript length属性
- javascript 字符串函数
- Javascript 日期对象简介
- Javascript 日期对象用途
- Date 对象属性和方法
- Javascript 数组是什么
- Javascript 创建数组
- Javascript 数组赋值与取值
- Javascript 数组属性和方法
- 实例解析php的数据类型
- 实现PHP中session存储及删除变量
- php微信公众号开发之秒杀
- php fread函数使用方法总结
- Yii2框架控制器、路由、Url生成操作示例
- Laravel框架实现调用百度翻译API功能示例
- phpstudy2018升级MySQL5.5为5.7教程(图文)
- laravel实现简单用户权限的示例代码
- tp5(thinkPHP5框架)时间查询操作实例分析
- PHP使Laravel为JSON REST API返回自定义错误的问题
- 详解PHP PDO简单教程
- Python实现ElGamal加密算法的示例代码
- PHP实现基于状态的责任链审批模式详解
- django rest framework使用django-filter用法
- 通过实例解析python创建进程常用方法