zoukankan      html  css  js  c++  java
  • Kafka.net使用编程入门(二)

    1.首先创建一个Topic,命令如下:

    kafka-topics --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic MyTopic

    2.创建两个控制台程序:

    3.KafkaProducer程序: 

    复制代码
    class Program
        {
            static void Main(string[] args)
            {
                do
                {
                    Produce(GetKafkaBroker(), getTopicName());
                    System.Threading.Thread.Sleep(3000);
                } while (true);
            }
    
            private static void Produce(string broker, string topic)
            {
                var options = new KafkaOptions(new Uri(broker));
                var router = new BrokerRouter(options);
                var client = new Producer(router);
    
                var currentDatetime =DateTime.Now;
                var key = currentDatetime.Second.ToString();
                var events = new[] { new Message("Hello World " + currentDatetime, key) };
                client.SendMessageAsync(topic, events).Wait(1500);
                Console.WriteLine("Produced: Key: {0}. Message: {1}", key, events[0].Value.ToUtf8String());
    
                using (client) { }
            }
    
            private static string GetKafkaBroker()
            {
                string KafkaBroker = string.Empty;
                const string kafkaBrokerKeyName = "KafkaBroker";
    
                if (!ConfigurationManager.AppSettings.AllKeys.Contains(kafkaBrokerKeyName))
                {
                    KafkaBroker = "http://localhost:9092";
                }
                else
                {
                    KafkaBroker = ConfigurationManager.AppSettings[kafkaBrokerKeyName];
                }
                return KafkaBroker;
            }
            private static string getTopicName()
            {
                string TopicName = string.Empty;
                const string topicNameKeyName = "Topic";
    
                if (!ConfigurationManager.AppSettings.AllKeys.Contains(topicNameKeyName))
                {
                    throw new Exception("Key "" + topicNameKeyName + "" not found in Config file -> configuration/AppSettings");
                }
                else
                {
                    TopicName = ConfigurationManager.AppSettings[topicNameKeyName];
                }
                return TopicName;
            }
        }
    复制代码

    4.KafkaConsumer程序:

    复制代码
    class Program
        {
            static void Main(string[] args)
            {
                Consume(getKafkaBroker(), getTopicName());
                
            }
    
            private static void Consume(string broker, string topic)
            {   
                var options = new KafkaOptions(new Uri(broker));
                var router = new BrokerRouter(options);
                var consumer = new Consumer(new ConsumerOptions(topic, router));
    
                //Consume returns a blocking IEnumerable (ie: never ending stream)
                foreach (var message in consumer.Consume())
                {
                    Console.WriteLine("Response: Partition {0},Offset {1} : {2}",
                        message.Meta.PartitionId, message.Meta.Offset, message.Value.ToUtf8String());
                }
            }
    
            private static string getKafkaBroker()
            {
                string KafkaBroker = string.Empty;
                var KafkaBrokerKeyName = "KafkaBroker";
    
                if (!ConfigurationManager.AppSettings.AllKeys.Contains(KafkaBrokerKeyName))
                {
                    KafkaBroker = "http://localhost:9092";
                }
                else
                {
                    KafkaBroker = ConfigurationManager.AppSettings[KafkaBrokerKeyName];
                }
                return KafkaBroker;
            }
    
            private static string getTopicName()
            {
                string TopicName = string.Empty;
                var TopicNameKeyName = "Topic";
    
                if (!ConfigurationManager.AppSettings.AllKeys.Contains(TopicNameKeyName))
                {
                    throw new Exception("Key "" + TopicNameKeyName + "" not found in Config file -> configuration/AppSettings");
                }
                else
                {
                    TopicName = ConfigurationManager.AppSettings[TopicNameKeyName];
                }
                return TopicName;
            }
        }
    复制代码

    5.Consumer结果:

  • 相关阅读:
    combo参数配置_手册
    mysql服务器辅助选项
    CentOS中操作
    Linux PHP增加JSON支持及如何使用JSON
    linux服务器命令
    linux中的工具
    linux文件夹操作(及模糊搜索)
    治疗肾结石
    其他书籍
    如何定位到div滚动条的最底端
  • 原文地址:https://www.cnblogs.com/zxtceq/p/9067831.html
Copyright © 2011-2022 走看看