怎么在Java中利用kafka发送消息-创新互联

这期内容当中小编将会给大家带来有关怎么在Java中利用kafka发送消息,文章内容丰富且以专业的角度为大家分析和叙述,阅读完这篇文章希望大家可以有所收获。

成都创新互联公司主营开平网站建设的网络公司,主营网站建设方案,重庆APP软件开发,开平h5成都微信小程序搭建,开平网站营销推广欢迎开平等地区企业咨询

1. maven依赖包

<dependency> 
 <groupId>org.apache.kafka</groupId> 
 <artifactId>kafka-clients</artifactId> 
 <version>0.9.0.1</version> 
</dependency>

2. 生产者代码

package com.lnho.example.kafka;  
import org.apache.kafka.clients.producer.KafkaProducer; 
import org.apache.kafka.clients.producer.Producer; 
import org.apache.kafka.clients.producer.ProducerRecord;   
import java.util.Properties;   
public class KafkaProducerExample { 
 public static void main(String[] args) { 
  Properties props = new Properties(); 
  props.put("bootstrap.servers", "master:9092"); 
  props.put("acks", "all"); 
  props.put("retries", 0); 
  props.put("batch.size", 16384); 
  props.put("linger.ms", 1); 
  props.put("buffer.memory", 33554432); 
  props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); 
  props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");   
  Producer<String, String> producer = new KafkaProducer<>(props); 
  for(int i = 0; i < 100; i++) 
   producer.send(new ProducerRecord<>("topic1", Integer.toString(i), Integer.toString(i)));   
  producer.close(); 
 } 
}

3. 消费者代码

package com.lnho.example.kafka;   
import org.apache.kafka.clients.consumer.ConsumerRecord; 
import org.apache.kafka.clients.consumer.ConsumerRecords; 
import org.apache.kafka.clients.consumer.KafkaConsumer; 
import java.util.Arrays; 
import java.util.Properties;   
public class KafkaConsumerExample { 
 public static void main(String[] args) { 
  Properties props = new Properties(); 
  props.put("bootstrap.servers", "master:9092"); 
  props.put("group.id", "test"); 
  props.put("enable.auto.commit", "true"); 
  props.put("auto.commit.interval.ms", "1000"); 
  props.put("session.timeout.ms", "30000"); 
  props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); 
  props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); 
  KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); 
  consumer.subscribe(Arrays.asList("topic1")); 
  while (true) { 
   ConsumerRecords<String, String> records = consumer.poll(100); 
   for (ConsumerRecord<String, String> record : records) 
    System.out.printf("offset = %d, key = %s, value = %s\n", record.offset(), record.key(), record.value()); 
  } 
 } 
}

上述就是小编为大家分享的怎么在Java中利用kafka发送消息了,如果刚好有类似的疑惑,不妨参照上述分析进行理解。如果想知道更多相关知识,欢迎关注创新互联行业资讯频道。

网页题目:怎么在Java中利用kafka发送消息-创新互联
标题链接:https://www.cdcxhl.com/article36/cdddpg.html

成都网站建设公司_创新互联,为您提供网站排名营销型网站建设云服务器用户体验标签优化企业建站

广告

声明:本网站发布的内容(图片、视频和文字)以用户投稿、用户转载内容为主,如果涉及侵权请尽快告知,我们将会在第一时间删除。文章观点不代表本网站立场,如需处理请联系客服。电话:028-86922220;邮箱:631063699@qq.com。内容未经允许不得转载,或转载时需注明来源: 创新互联

营销型网站建设