我的编程空间,编程开发者的网络收藏夹
学习永远不晚

大数据之Kafka————java来实现kafka相关操作

短信预约 -IT技能 免费直播动态提醒
省份

北京

  • 北京
  • 上海
  • 天津
  • 重庆
  • 河北
  • 山东
  • 辽宁
  • 黑龙江
  • 吉林
  • 甘肃
  • 青海
  • 河南
  • 江苏
  • 湖北
  • 湖南
  • 江西
  • 浙江
  • 广东
  • 云南
  • 福建
  • 海南
  • 山西
  • 四川
  • 陕西
  • 贵州
  • 安徽
  • 广西
  • 内蒙
  • 西藏
  • 新疆
  • 宁夏
  • 兵团
手机号立即预约

请填写图片验证码后获取短信验证码

看不清楚,换张图片

免费获取短信验证码

大数据之Kafka————java来实现kafka相关操作

一、在java中配置pom

           junit      junit      4.11      test              org.apache.kafka      kafka-clients      2.8.0              org.apache.kafka      kafka_2.12      2.8.0      

二、生产者方法

(1)、Producer

Java中写在生产者输入内容在kafka中可以让消费者提取

[root@kb144 config]# kafka-console-consumer.sh --bootstrap-server 192.168.153.144:9092 --topic kb22

package nj.zb.kb22.Kafka;import org.apache.kafka.clients.producer.KafkaProducer;import org.apache.kafka.clients.producer.ProducerConfig;import org.apache.kafka.clients.producer.ProducerRecord;import org.apache.kafka.common.serialization.StringSerializer;import java.util.Properties;import java.util.Scanner;public class MyProducer {    public static void main(String[] args) {        Properties properties = new Properties();        //生产者的配置文件        properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"192.168.153.144:9092");        //key的序列化        properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);        //value的序列化        properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,StringSerializer.class);               properties.put(ProducerConfig.ACKS_CONFIG,"1");        KafkaProducer producer = new KafkaProducer(properties);        Scanner scanner = new Scanner(System.in);        while (true){            System.out.println("请输入kafka的内容");            String msg =scanner.next();            ProducerRecord record = new ProducerRecord("kb22",msg);            producer.send(record);        }    }}

(2)、Producer进行多线程操作

  生产者多线程是一种常见的技术实践,可以提高消息生产的并发性和吞吐量。通过将消息生产任务分配给多个线程来并行地发送消息,可以有效地利用系统资源,加快消息的发送速度。

package nj.zb.kb22.Kafka;import org.apache.kafka.clients.producer.KafkaProducer;import org.apache.kafka.clients.producer.ProducerConfig;import org.apache.kafka.clients.producer.ProducerRecord;import org.apache.kafka.common.serialization.StringSerializer;import java.util.Properties;import java.util.concurrent.ExecutorService;import java.util.concurrent.Executors;public class MyProducer2 {    public static void main(String[] args) {        ExecutorService executorService = Executors.newCachedThreadPool();        for (int i = 0; i < 10; i++) {//i代表线程            Thread thread =new Thread(new Runnable() {                @Override                public void run() {                    Properties properties = new Properties();                      properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"192.168.153.144:9092");   properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);                properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,StringSerializer.class);  properties.put(ProducerConfig.ACKS_CONFIG,"0");  KafkaProducer producer = new KafkaProducer(properties);                    //多线程操作 j代表消息                    for (int j = 0; j < 100; j++) {                        String msg=Thread.currentThread().getName()+" "+ j;                        System.out.println(msg);                        ProducerRecord re = new ProducerRecord("kb22", msg);                        producer.send(re);                    }                }            });            executorService.execute(thread);        }        executorService.shutdown();        while (true){            if (executorService.isTerminated()){                System.out.println("game over");                break;            }        }    }}

三、消费者方法

(1)、Consumer

通过java来实现消费者

package nj.zb.kb22.Kafka;import org.apache.kafka.clients.consumer.ConsumerConfig;import org.apache.kafka.clients.consumer.ConsumerRecord;import org.apache.kafka.clients.consumer.ConsumerRecords;import org.apache.kafka.clients.consumer.KafkaConsumer;import org.apache.kafka.common.serialization.StringDeserializer;import org.apache.kafka.common.serialization.StringSerializer;import java.time.Duration;import java.util.Collections;import java.util.Properties;public class MyConsumer {    public static void main(String[] args) {        Properties properties = new Properties();        properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,"192.168.153.144:9092");        properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);        properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class);        //设置拉取信息后是否自动提交(kafka记录当前app是否已经获取到此信息),false 手动提交 ;true 自动提交        properties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,"false");                properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,"earliest");        properties.put(ConsumerConfig.GROUP_ID_CONFIG,"GROUP3");        KafkaConsumer consumer = new KafkaConsumer(properties);        //创建好kafka消费者对象后,订阅消息,指定消费的topic        consumer.subscribe(Collections.singleton("kb22"));        while (true){            Duration mills = Duration.ofMillis(100);            ConsumerRecords records = consumer.poll(mills);            for (ConsumerRecord record:records){                String topic = record.topic();                int partition = record.partition();                long offset = record.offset();                String key = record.key();                String value = record.value();                long timestamp = record.timestamp();                System.out.println("topic:"+topic+"\tpartition"+partition+"\toffset"+offset+"\tkey"+key+"\tvalue"+value+"\ttimestamp"+timestamp);            }            //consumer.commitAsync();//手动提交        }    }}

(2)、设置多人访问

package nj.zb.kb22.Kafka;import org.apache.kafka.clients.consumer.ConsumerConfig;import org.apache.kafka.clients.consumer.ConsumerRecord;import org.apache.kafka.clients.consumer.ConsumerRecords;import org.apache.kafka.clients.consumer.KafkaConsumer;import org.apache.kafka.common.serialization.StringDeserializer;import java.time.Duration;import java.util.Collections;import java.util.Properties;public class MyConsumerThread {    //模仿多人访问    public static void main(String[] args) {        for (int i = 0; i <3; i++) {            new Thread(new Runnable() {                @Override                public void run() {                    Properties properties = new Properties();                    properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,"192.168.153.144:9092");                    properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);                    properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class);                    //设置拉取信息后是否自动提交(kafka记录当前app是否已经获取到此信息),false 手动提交 ;true 自动提交                    properties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,"false");                                        properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,"earliest");                    properties.put(ConsumerConfig.GROUP_ID_CONFIG,"GROUP3");                    KafkaConsumer consumer = new KafkaConsumer<>(properties);                    consumer.subscribe(Collections.singleton("kb22"));                    while (true){                        //poll探寻数据                        ConsumerRecords records = consumer.poll(Duration.ofMillis(100));                        for (ConsumerRecordrecord:records){String topic = record.topic();int partition = record.partition();long offset = record.offset();String key = record.key();String value = record.value();long timestamp = record.timestamp();String name = Thread.currentThread().getName();System.out.println("name"+name        +"\ttopic:"+topic        +"\tpartition" +partition        +"\toffset"+offset        +"\tkey"+key        +"\tvalue"+value        +"\ttimestamp"+timestamp);                        }                    }                }            }).start();        }    }}

来源地址:https://blog.csdn.net/ycz926940/article/details/131562785

免责声明:

① 本站未注明“稿件来源”的信息均来自网络整理。其文字、图片和音视频稿件的所属权归原作者所有。本站收集整理出于非商业性的教育和科研之目的,并不意味着本站赞同其观点或证实其内容的真实性。仅作为临时的测试数据,供内部测试之用。本站并未授权任何人以任何方式主动获取本站任何信息。

② 本站未注明“稿件来源”的临时测试数据将在测试完成后最终做删除处理。有问题或投稿请发送至: 邮箱/279061341@qq.com QQ/279061341

大数据之Kafka————java来实现kafka相关操作

下载Word文档到电脑,方便收藏和打印~

下载Word文档

编程热搜

  • Python 学习之路 - Python
    一、安装Python34Windows在Python官网(https://www.python.org/downloads/)下载安装包并安装。Python的默认安装路径是:C:\Python34配置环境变量:【右键计算机】--》【属性】-
    Python 学习之路 - Python
  • chatgpt的中文全称是什么
    chatgpt的中文全称是生成型预训练变换模型。ChatGPT是什么ChatGPT是美国人工智能研究实验室OpenAI开发的一种全新聊天机器人模型,它能够通过学习和理解人类的语言来进行对话,还能根据聊天的上下文进行互动,并协助人类完成一系列
    chatgpt的中文全称是什么
  • C/C++中extern函数使用详解
  • C/C++可变参数的使用
    可变参数的使用方法远远不止以下几种,不过在C,C++中使用可变参数时要小心,在使用printf()等函数时传入的参数个数一定不能比前面的格式化字符串中的’%’符号个数少,否则会产生访问越界,运气不好的话还会导致程序崩溃
    C/C++可变参数的使用
  • css样式文件该放在哪里
  • php中数组下标必须是连续的吗
  • Python 3 教程
    Python 3 教程 Python 的 3.0 版本,常被称为 Python 3000,或简称 Py3k。相对于 Python 的早期版本,这是一个较大的升级。为了不带入过多的累赘,Python 3.0 在设计的时候没有考虑向下兼容。 Python
    Python 3 教程
  • Python pip包管理
    一、前言    在Python中, 安装第三方模块是通过 setuptools 这个工具完成的。 Python有两个封装了 setuptools的包管理工具: easy_install  和  pip , 目前官方推荐使用 pip。    
    Python pip包管理
  • ubuntu如何重新编译内核
  • 改善Java代码之慎用java动态编译

目录