日韩性视频-久久久蜜桃-www中文字幕-在线中文字幕av-亚洲欧美一区二区三区四区-撸久久-香蕉视频一区-久久无码精品丰满人妻-国产高潮av-激情福利社-日韩av网址大全-国产精品久久999-日本五十路在线-性欧美在线-久久99精品波多结衣一区-男女午夜免费视频-黑人极品ⅴideos精品欧美棵-人人妻人人澡人人爽精品欧美一区-日韩一区在线看-欧美a级在线免费观看

歡迎訪問 生活随笔!

生活随笔

當(dāng)前位置: 首頁 > 编程资源 > 编程问答 >内容正文

编程问答

java利用kafka生产消费消息

發(fā)布時(shí)間:2025/4/5 编程问答 28 豆豆
生活随笔 收集整理的這篇文章主要介紹了 java利用kafka生产消费消息 小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,幫大家做個(gè)參考.

2019獨(dú)角獸企業(yè)重金招聘Python工程師標(biāo)準(zhǔn)>>>

1.producer程序

package com.test.frame.kafka.controller;import kafka.javaapi.producer.Producer; import kafka.producer.KeyedMessage; import kafka.producer.ProducerConfig;import java.util.Properties;public class KafkaProducer {private final Producer<String, String> producer;public final static String TOPIC = "my-multi-topic";//構(gòu)造方法private KafkaProducer() {Properties props = new Properties();props.put("metadata.broker.list", "localhost:9092");props.put("serializer.class", "kafka.serializer.StringEncoder");props.put("key.serializer.class", "kafka.serializer.StringEncoder");props.put("request.required.acks", "-1");producer = new Producer<String, String>(new ProducerConfig(props));}void produce() {int messageNo = 90;final int COUNT = 100;while (messageNo < COUNT) {String key = String.valueOf(messageNo);String data = "hello kafka message" + key;producer.send(new KeyedMessage<String, String>(TOPIC, key ,data));System.out.println(data);messageNo++;}}public static void main(String[] args) throws Exception {new KafkaProducer().produce();}}

運(yùn)行結(jié)果:

消費(fèi)方接收到的消息如下:

2.consumer端程序:

package com.test.frame.kafka.controller;import kafka.consumer.ConsumerConfig; import kafka.consumer.ConsumerIterator; import kafka.consumer.KafkaStream; import kafka.javaapi.consumer.ConsumerConnector; import kafka.serializer.StringDecoder; import kafka.utils.VerifiableProperties;import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Properties;public class KafkaConsumer {private final ConsumerConnector consumer;private KafkaConsumer() {Properties props = new Properties();//zookeeper 配置props.put("zookeeper.connect", "localhost:2181");//group 代表一個(gè)消費(fèi)組props.put("group.id", "jd-group");//zk連接超時(shí)props.put("zookeeper.session.timeout.ms", "4000");props.put("zookeeper.sync.time.ms", "200");props.put("auto.commit.interval.ms", "1000");props.put("auto.offset.reset", "smallest");//序列化類props.put("serializer.class", "kafka.serializer.StringEncoder");ConsumerConfig config = new ConsumerConfig(props);consumer = kafka.consumer.Consumer.createJavaConsumerConnector(config); }void consume() {Map<String, Integer> topicCountMap = new HashMap<String, Integer>();topicCountMap.put(KafkaProducer.TOPIC, new Integer(1));StringDecoder keyDecoder = new StringDecoder(new VerifiableProperties());StringDecoder valueDecoder = new StringDecoder(new VerifiableProperties());Map<String, List<KafkaStream<String, String>>> consumerMap =consumer.createMessageStreams(topicCountMap,keyDecoder,valueDecoder);KafkaStream<String, String> stream = consumerMap.get(KafkaProducer.TOPIC).get(0);ConsumerIterator<String, String> it = stream.iterator();while (it.hasNext())System.out.println(it.next().message());}public static void main(String[] args) {new KafkaConsumer().consume();}}

運(yùn)行結(jié)果如下:

此時(shí)已經(jīng)聯(lián)通成功。

?

?

?

?

?

轉(zhuǎn)載于:https://my.oschina.net/u/2263272/blog/1527979

總結(jié)

以上是生活随笔為你收集整理的java利用kafka生产消费消息的全部內(nèi)容,希望文章能夠幫你解決所遇到的問題。

如果覺得生活随笔網(wǎng)站內(nèi)容還不錯(cuò),歡迎將生活随笔推薦給好友。