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

Kafka Java客户端代码的示例分析

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

北京

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

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

看不清楚,换张图片

免费获取短信验证码

Kafka Java客户端代码的示例分析

这篇文章将为大家详细讲解有关Kafka Java客户端代码的示例分析,文章内容质量较高,因此小编分享给大家做个参考,希望大家阅读完这篇文章后对相关知识有一定的了解。

kafka是一种高吞吐量的分布式发布订阅消息系统

kafka是linkedin用于日志处理的分布式消息队列,linkedin的日志数据容量大,但对可靠性要求不高,其日志数据主要包括用户行为(登录、浏览、点击、分享、喜欢)以及系统运行日志(CPU、内存、磁盘、网络、系统及进程状态)

当前很多的消息队列服务提供可靠交付保证,并默认是即时消费(不适合离线)。

高可靠交付对linkedin的日志不是必须的,故可通过降低可靠性来提高性能,同时通过构建分布式的集群,允许消息在系统中累积,使得kafka同时支持离线和在线日志处理

测试环境

kafka_2.10-0.8.1.1 3个节点做的集群

zookeeper-3.4.5 一个实例节点

代码示例

消息生产者代码示例

import java.util.Collections;  import java.util.Date;  import java.util.Properties;  import java.util.Random;     import kafka.javaapi.producer.Producer;  import kafka.producer.KeyedMessage;  import kafka.producer.ProducerConfig;      public class ProducerDemo {      public static void main(String[] args) {          Random rnd = new Random();          int events=100;             // 设置配置属性          Properties props = new Properties();          props.put("metadata.broker.list","172.168.63.221:9092,172.168.63.233:9092,172.168.63.234:9092");          props.put("serializer.class", "kafka.serializer.StringEncoder");          // key.serializer.class默认为serializer.class          props.put("key.serializer.class", "kafka.serializer.StringEncoder");          // 可选配置,如果不配置,则使用默认的partitioner          props.put("partitioner.class", "com.catt.kafka.demo.PartitionerDemo");          // 触发acknowledgement机制,否则是fire and forget,可能会引起数据丢失          // 值为0,1,-1,可以参考          // http://kafka.apache.org/08/configuration.html          props.put("request.required.acks", "1");          ProducerConfig config = new ProducerConfig(props);             // 创建producer          Producer<String, String> producer = new Producer<String, String>(config);          // 产生并发送消息          long start=System.currentTimeMillis();          for (long i = 0; i < events; i++) {              long runtime = new Date().getTime();              String ip = "192.168.2." + i;//rnd.nextInt(255);              String msg = runtime + ",www.example.com," + ip;              //如果topic不存在,则会自动创建,默认replication-factor为1,partitions为0              KeyedMessage<String, String> data = new KeyedMessage<String, String>(                      "page_visits", ip, msg);              producer.send(data);          }          System.out.println("耗时:" + (System.currentTimeMillis() - start));          // 关闭producer          producer.close();      }  }

消息消费者代码示例

import java.util.HashMap;  import java.util.List;  import java.util.Map;  import java.util.Properties;  import java.util.concurrent.ExecutorService;  import java.util.concurrent.Executors;     import kafka.consumer.Consumer;  import kafka.consumer.ConsumerConfig;  import kafka.consumer.KafkaStream;  import kafka.javaapi.consumer.ConsumerConnector;      public class ConsumerDemo {      private final ConsumerConnector consumer;      private final String topic;      private ExecutorService executor;         public ConsumerDemo(String a_zookeeper, String a_groupId, String a_topic) {          consumer = Consumer.createJavaConsumerConnector(createConsumerConfig(a_zookeeper,a_groupId));          this.topic = a_topic;      }         public void shutdown() {          if (consumer != null)              consumer.shutdown();          if (executor != null)              executor.shutdown();      }         public void run(int numThreads) {          Map<String, Integer> topicCountMap = new HashMap<String, Integer>();          topicCountMap.put(topic, new Integer(numThreads));          Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = consumer                  .createMessageStreams(topicCountMap);          List<KafkaStream<byte[], byte[]>> streams = consumerMap.get(topic);             // now launch all the threads          executor = Executors.newFixedThreadPool(numThreads);             // now create an object to consume the messages          //          int threadNumber = 0;          for (final KafkaStream stream : streams) {              executor.submit(new ConsumerMsgTask(stream, threadNumber));              threadNumber++;          }      }         private static ConsumerConfig createConsumerConfig(String a_zookeeper,              String a_groupId) {          Properties props = new Properties();          props.put("zookeeper.connect", a_zookeeper);          props.put("group.id", a_groupId);          props.put("zookeeper.session.timeout.ms", "400");          props.put("zookeeper.sync.time.ms", "200");          props.put("auto.commit.interval.ms", "1000");             return new ConsumerConfig(props);      }         public static void main(String[] arg) {          String[] args = { "172.168.63.221:2188", "group-1", "page_visits", "12" };          String zooKeeper = args[0];          String groupId = args[1];          String topic = args[2];          int threads = Integer.parseInt(args[3]);             ConsumerDemo demo = new ConsumerDemo(zooKeeper, groupId, topic);          demo.run(threads);             try {              Thread.sleep(10000);          } catch (InterruptedException ie) {             }          demo.shutdown();      }  }

消息处理类

import kafka.consumer.ConsumerIterator;  import kafka.consumer.KafkaStream;     public class ConsumerMsgTask implements Runnable {      private KafkaStream m_stream;      private int m_threadNumber;         public ConsumerMsgTask(KafkaStream stream, int threadNumber) {          m_threadNumber = threadNumber;          m_stream = stream;      }         public void run() {          ConsumerIterator<byte[], byte[]> it = m_stream.iterator();          while (it.hasNext())              System.out.println("Thread " + m_threadNumber + ": "                     + new String(it.next().message()));          System.out.println("Shutting down Thread: " + m_threadNumber);      }  }

Partitioner类示例

import kafka.producer.Partitioner;  import kafka.utils.VerifiableProperties;     public class PartitionerDemo implements Partitioner {      public PartitionerDemo(VerifiableProperties props) {         }         @Override     public int partition(Object obj, int numPartitions) {          int partition = 0;          if (obj instanceof String) {              String key=(String)obj;              int offset = key.lastIndexOf('.');              if (offset > 0) {                  partition = Integer.parseInt(key.substring(offset + 1)) % numPartitions;              }          }else{              partition = obj.toString().length() % numPartitions;          }                     return partition;      }     }

关于Kafka Java客户端代码的示例分析就分享到这里了,希望以上内容可以对大家有一定的帮助,可以学到更多知识。如果觉得文章不错,可以把它分享出去让更多的人看到。

免责声明:

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

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

Kafka Java客户端代码的示例分析

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

下载Word文档

猜你喜欢

Kafka Java客户端代码的示例分析

这篇文章将为大家详细讲解有关Kafka Java客户端代码的示例分析,文章内容质量较高,因此小编分享给大家做个参考,希望大家阅读完这篇文章后对相关知识有一定的了解。kafka是一种高吞吐量的分布式发布订阅消息系统kafka是linkedin
2023-06-17

Kafka简单客户端编程的示例分析

这篇文章将为大家详细讲解有关Kafka简单客户端编程的示例分析,小编觉得挺实用的,因此分享给大家做个参考,希望大家阅读完这篇文章后可以有所收获。一、创建配置类Config这个类很简单,只是存放了两个常量,一个是话题TOPIC,一个是线程数T
2023-05-30

Golang语言HTTP客户端的示例分析

这篇文章将为大家详细讲解有关Golang语言HTTP客户端的示例分析,小编觉得挺实用的,因此分享给大家做个参考,希望大家阅读完这篇文章后可以有所收获。HTTP客户端封装package taskimport ( "bytes" "encodi
2023-06-25

Java中http下载文件客户端和上传文件客户端的示例分析

这篇文章主要介绍了Java中http下载文件客户端和上传文件客户端的示例分析,具有一定借鉴价值,感兴趣的朋友可以参考下,希望大家阅读完这篇文章之后大有收获,下面让小编带着大家一起了解一下。一、下载客户端代码package javadownl
2023-05-30

Spring-Boot集成Solr客户端的示例分析

这篇文章主要为大家展示了“Spring-Boot集成Solr客户端的示例分析”,内容简而易懂,条理清晰,希望能够帮助大家解决疑惑,下面让小编带领大家一起研究并学习一下“Spring-Boot集成Solr客户端的示例分析”这篇文章吧。Solr
2023-05-30

Java客户端利用Jedis操作redis缓存示例代码

前言Redis是一个开源的Key-Value数据缓存,和Memcached类似。Redis多种类型的value,包括string(字符串)、list(链表)、set(集合)、zset(sorted set --有序集合)和hash(哈希类型
2023-05-31

C# winfroms使用socket客户端服务端的示例代码

本文介绍了如何在C#WinForms中使用Socket实现客户端-服务器通信。示例代码演示了客户端如何连接到服务器、发送和接收数据,以及服务器如何接收连接、发送和接收数据。本文适用于希望在C#WinForms应用程序中实现网络通信的开发人员。
C# winfroms使用socket客户端服务端的示例代码
2024-04-02

客户端JavaScript线程池设计的示例分析

这篇“客户端JavaScript线程池设计的示例分析”除了程序员外大部分人都不太理解,今天小编为了让大家更加理解“客户端JavaScript线程池设计的示例分析”,给大家总结了以下内容,具有一定借鉴价值,内容详细步骤清晰,细节处理妥当,希望
2023-06-28

vue3+electron12+dll开发客户端配置的示例分析

小编给大家分享一下vue3+electron12+dll开发客户端配置的示例分析,相信大部分人都还不怎么了解,因此分享这篇文章给大家参考一下,希望大家阅读完这篇文章后大有收获,下面让我们一起去了解一下吧!修改仓库源由于electron版本的
2023-06-15

如何进行Kafka 1.0.0 d代码示例分析

这篇文章将为大家详细讲解有关如何进行Kafka 1.0.0 d代码示例分析,文章内容质量较高,因此小编分享给大家做个参考,希望大家阅读完这篇文章后对相关知识有一定的了解。package kafka.demo;import java.util
2023-06-02

Java数组代码的示例分析

本篇文章给大家分享的是有关Java数组代码的示例分析,小编觉得挺实用的,因此分享给大家学习,希望大家阅读完这篇文章后可以有所收获,话不多说,跟着小编一起来看看吧。数组分类1. 一维数组1.1 一维数组的定义和初始化1.2 对一维数组的操作,
2023-06-02

UDP服务器客户端编程流程的示例分析

这篇文章给大家分享的是有关UDP服务器客户端编程流程的示例分析的内容。小编觉得挺实用的,因此分享给大家做个参考,一起跟随小编过来看看吧。UDP编程流程UDP提供的是无连接、不可靠的、数据报服务UDP是尽最大能力进行传输,但是并不能保证可靠性
2023-06-21

Netty分布式客户端处理接入事件handle的示例分析

这篇文章主要介绍了Netty分布式客户端处理接入事件handle的示例分析,具有一定借鉴价值,感兴趣的朋友可以参考下,希望大家阅读完这篇文章之后大有收获,下面让小编带着大家一起了解一下。处理接入事件创建handle回到上一章NioEvent
2023-06-29

编程热搜

  • 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动态编译

目录