热门标签 | HotTags
当前位置:  开发笔记 > 编程语言 > 正文

深入解析Java如何利用Redis实现高效消息队列

本篇文章主要介绍了Java利用Redis实现消息队列的示例代码,小编觉得挺不错的,现在分享给大家,也给大家做个参考。一起跟随小编过来看看吧本文介绍了Java利用Redis实现消息队

本篇文章主要介绍了Java利用Redis实现消息队列的示例代码,小编觉得挺不错的,现在分享给大家,也给大家做个参考。一起跟随小编过来看看吧

本文介绍了Java利用Redis实现消息队列的示例代码,分享给大家,具体如下:

应用场景

为什么要用redis?

二进制存储、java序列化传输、IO连接数高、连接频繁

一、序列化

这里编写了一个java序列化的工具,主要是将对象转化为byte数组,和根据byte数组反序列化成java对象; 主要是用到了ByteArrayOutputStream和ByteArrayInputStream; 注意:每个需要序列化的对象都要实现Serializable接口;

其代码如下:

package Utils;import java.io.*;/** * Created by Kinglf on 2016/10/17. */public class ObjectUtil { /**  * 对象转byte[]  * @param obj  * @return  * @throws IOException  */ public static byte[] object2Bytes(Object obj) throws IOException{  ByteArrayOutputStream bo=new ByteArrayOutputStream();  ObjectOutputStream oo=new ObjectOutputStream(bo);  oo.writeObject(obj);  byte[] bytes=bo.toByteArray();  bo.close();  oo.close();  return bytes; } /**  * byte[]转对象  * @param bytes  * @return  * @throws Exception  */ public static Object bytes2Object(byte[] bytes) throws Exception{  ByteArrayInputStream in=new ByteArrayInputStream(bytes);  ObjectInputStream sIn=new ObjectInputStream(in);  return sIn.readObject(); }}

二、消息类(实现Serializable接口)

package Model;import java.io.Serializable;/** * Created by Kinglf on 2016/10/17. */public class Message implements Serializable { private static final long serialVersiOnUID= -389326121047047723L; private int id; private String content; public Message(int id, String content) {  this.id = id;  this.cOntent= content; } public int getId() {  return id; } public void setId(int id) {  this.id = id; } public String getContent() {  return content; } public void setContent(String content) {  this.cOntent= content; }}

三、Redis的操作

利用redis做队列,我们采用的是redis中list的push和pop操作;

结合队列的特点:

只允许在一端插入新元素只能在队列的尾部FIFO:先进先出原则 Redis中lpush头入(rpop尾出)或rpush尾入(lpop头出)可以满足要求,而Redis中list药push或 pop的对象仅需要转换成byte[]即可

java采用Jedis进行Redis的存储和Redis的连接池设置

上代码:

package Utils;import redis.clients.jedis.Jedis;import redis.clients.jedis.JedisPool;import redis.clients.jedis.JedisPoolConfig;import java.util.List;import java.util.Map;import java.util.Set;/** * Created by Kinglf on 2016/10/17. */public class JedisUtil { private static String JEDIS_IP; private static int JEDIS_PORT; private static String JEDIS_PASSWORD; private static JedisPool jedisPool; static {  //Configuration自行写的配置文件解析类,继承自Properties  Configuration cOnf=Configuration.getInstance();  JEDIS_IP=conf.getString("jedis.ip","127.0.0.1");  JEDIS_PORT=conf.getInt("jedis.port",6379);  JEDIS_PASSWORD=conf.getString("jedis.password",null);  JedisPoolConfig cOnfig=new JedisPoolConfig();  config.setMaxActive(5000);  config.setMaxIdle(256);  config.setMaxWait(5000L);  config.setTestOnBorrow(true);  config.setTestOnReturn(true);  config.setTestWhileIdle(true);  config.setMinEvictableIdleTimeMillis(60000L);  config.setTimeBetweenEvictionRunsMillis(3000L);  config.setNumTestsPerEvictionRun(-1);  jedisPool=new JedisPool(config,JEDIS_IP,JEDIS_PORT,60000); } /**  * 获取数据  * @param key  * @return  */ public static String get(String key){  String value=null;  Jedis jedis=null;  try{   jedis=jedisPool.getResource();   value=jedis.get(key);  }catch (Exception e){   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  }finally {   close(jedis);  }  return value; } private static void close(Jedis jedis) {  try{   jedisPool.returnResource(jedis);  }catch (Exception e){   if(jedis.isConnected()){    jedis.quit();    jedis.disconnect();   }  } } public static byte[] get(byte[] key){  byte[] value = null;  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   value = jedis.get(key);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  }  return value; } public static void set(byte[] key, byte[] value) {  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.set(key, value);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  } } public static void set(byte[] key, byte[] value, int time) {  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.set(key, value);   jedis.expire(key, time);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  } } public static void hset(byte[] key, byte[] field, byte[] value) {  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.hset(key, field, value);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  } } public static void hset(String key, String field, String value) {  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.hset(key, field, value);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  } } /**  * 获取数据  *  * @param key  * @return  */ public static String hget(String key, String field) {  String value = null;  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   value = jedis.hget(key, field);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  }  return value; } /**  * 获取数据  *  * @param key  * @return  */ public static byte[] hget(byte[] key, byte[] field) {  byte[] value = null;  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   value = jedis.hget(key, field);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  }  return value; } public static void hdel(byte[] key, byte[] field) {  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.hdel(key, field);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  } } /**  * 存储REDIS队列 顺序存储  * @param key reids键名  * @param value 键值  */ public static void lpush(byte[] key, byte[] value) {  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.lpush(key, value);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  } } /**  * 存储REDIS队列 反向存储  * @param key reids键名  * @param value 键值  */ public static void rpush(byte[] key, byte[] value) {  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.rpush(key, value);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  } } /**  * 将列表 source 中的最后一个元素(尾元素)弹出,并返回给客户端  * @param key reids键名  * @param destination 键值  */ public static void rpoplpush(byte[] key, byte[] destination) {  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.rpoplpush(key, destination);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  } } /**  * 获取队列数据  * @param key 键名  * @return  */ public static List lpopList(byte[] key) {  List list = null;  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   list = jedis.lrange(key, 0, -1);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  }  return list; } /**  * 获取队列数据  * @param key 键名  * @return  */ public static byte[] rpop(byte[] key) {  byte[] bytes = null;  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   bytes = jedis.rpop(key);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  }  return bytes; } public static void hmset(Object key, Map hash) {  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.hmset(key.toString(), hash);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  } } public static void hmset(Object key, Map hash, int time) {  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.hmset(key.toString(), hash);   jedis.expire(key.toString(), time);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  } } public static List hmget(Object key, String... fields) {  List result = null;  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   result = jedis.hmget(key.toString(), fields);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  }  return result; } public static Set hkeys(String key) {  Set result = null;  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   result = jedis.hkeys(key);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  }  return result; } public static List lrange(byte[] key, int from, int to) {  List result = null;  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   result = jedis.lrange(key, from, to);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  }  return result; } public static Map hgetAll(byte[] key) {  Map result = null;  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   result = jedis.hgetAll(key);  } catch (Exception e) {   //释放redis对象   jed本文来源gaodai$ma#com搞$$代**码网isPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  }  return result; } public static void del(byte[] key) {  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.del(key);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  } } public static long llen(byte[] key) {  long len = 0;  Jedis jedis = null;  try {   jedis = jedisPool.getResource();   jedis.llen(key);  } catch (Exception e) {   //释放redis对象   jedisPool.returnBrokenResource(jedis);   e.printStackTrace();  } finally {   //返还到连接池   close(jedis);  }  return len; }}

四、Configuration主要用于读取Redis的配置信息

package Utils;import java.io.IOException;import java.io.InputStream;import java.util.Properties;/** * Created by Kinglf on 2016/10/17. */public class Configuration extends Properties { private static final long serialVersiOnUID= -2296275030489943706L; private static Configuration instance = null; public static synchronized Configuration getInstance() {  if (instance == null) {   instance = new Configuration();  }  return instance; } public String getProperty(String key, String defaultValue) {  String val = getProperty(key);  return (val == null || val.isEmpty()) ? defaultValue : val; } public String getString(String name, String defaultValue) {  return this.getProperty(name, defaultValue); } public int getInt(String name, int defaultValue) {  String val = this.getProperty(name);  return (val == null || val.isEmpty()) ? defaultValue : Integer.parseInt(val); } public long getLong(String name, long defaultValue) {  String val = this.getProperty(name);  return (val == null || val.isEmpty()) ? defaultValue : Integer.parseInt(val); } public float getFloat(String name, float defaultValue) {  String val = this.getProperty(name);  return (val == null || val.isEmpty()) ? defaultValue : Float.parseFloat(val); } public double getDouble(String name, double defaultValue) {  String val = this.getProperty(name);  return (val == null || val.isEmpty()) ? defaultValue : Double.parseDouble(val); } public byte getByte(String name, byte defaultValue) {  String val = this.getProperty(name);  return (val == null || val.isEmpty()) ? defaultValue : Byte.parseByte(val); } public Configuration() {  InputStream in = ClassLoader.getSystemClassLoader().getResourceAsStream("config.xml");  try {   this.loadFromXML(in);   in.close();  } catch (IOException ioe) {  } }}

五、测试

import Model.Message;import Utils.JedisUtil;import Utils.ObjectUtil;import redis.clients.jedis.Jedis;import java.io.IOException;/** * Created by Kinglf on 2016/10/17. */public class TestRedisQueue { public static byte[] redisKey = "key".getBytes(); static {  try {   init();  } catch (IOException e) {   e.printStackTrace();  } } private static void init() throws IOException {  for (int i = 0; i <1000000; i++) {   Message message = new Message(i, "这是第" + i + "个内容");   JedisUtil.lpush(redisKey, ObjectUtil.object2Bytes(message));  } } public static void main(String[] args) {  try {   pop();  } catch (Exception e) {   e.printStackTrace();  } } private static void pop() throws Exception {  byte[] bytes = JedisUtil.rpop(redisKey);  Message msg = (Message) ObjectUtil.bytes2Object(bytes);  if (msg != null) {   System.out.println(msg.getId() + "----" + msg.getContent());  } }}

每执行一次pop()方法,结果如下:
1----这是第1个内容
2----这是第2个内容
3----这是第3个内容
4----这是第4个内容

总结

至此,整个Redis消息队列的生产者和消费者代码已经完成

1.Message 需要传送的实体类(需实现Serializable接口)

2.Configuration Redis的配置读取类,继承自Properties

3.ObjectUtil 将对象和byte数组双向转换的工具类

4.Jedis 通过消息队列的先进先出(FIFO)的特点结合Redis的list中的push和pop操作进行封装的工具类

以上就是Java如何使用Redis来实现消息队列的具体分析的详细内容,更多请关注gaodaima其它相关文章!



推荐阅读
  • 本文介绍了如何使用Java编程语言实现凯撒密码的加密与解密功能。凯撒密码是一种替换式密码,通过将字母表中的每个字母向前或向后移动固定数量的位置来实现加密。 ... [详细]
  • 个人博客:打开链接依赖倒置原则定义依赖倒置原则(DependenceInversionPrinciple,DIP)定义如下:Highlevelmo ... [详细]
  • 使用Java计算两个日期之间的月份数
    本文详细介绍了利用Java编程语言计算两个指定日期之间月份数的方法。文章通过实例代码讲解了如何使用Joda-Time库来简化日期处理过程,旨在为开发者提供一个高效且易于理解的解决方案。 ... [详细]
  • Hadoop MapReduce 实战案例:手机流量使用统计分析
    本文通过一个具体的Hadoop MapReduce案例,详细介绍了如何利用MapReduce框架来统计和分析手机用户的流量使用情况,包括上行和下行流量的计算以及总流量的汇总。 ... [详细]
  • 本文基于Java官方文档进行了适当修改,旨在介绍如何实现一个能够同时处理多个客户端请求的服务端程序。在前文中,我们探讨了单客户端访问的服务端实现,而本篇将深入讲解多客户端环境下的服务端设计与实现。 ... [详细]
  • java datarow_DataSet  DataTable DataRow 深入浅出
    本篇文章适合有一定的基础的人去查看,最好学习过一定net编程基础在来查看此文章。1.概念DataSet是ADO.NET的中心概念。可以把DataSet当成内存中的数据 ... [详细]
  • 本文介绍如何通过Java代码调用阿里云短信服务API来实现短信验证码的发送功能,包括必要的依赖添加和关键代码示例。 ... [详细]
  • 使用 ModelAttribute 实现页面数据自动填充
    本文介绍了如何利用 Spring MVC 中的 ModelAttribute 注解,在页面跳转后自动填充表单数据。主要探讨了两种实现方法及其背后的原理。 ... [详细]
  • 本文介绍了一种在 Android 开发中动态修改 strings.xml 文件中字符串值的有效方法。通过使用占位符,开发者可以在运行时根据需要填充具体的值,从而提高应用的灵活性和可维护性。 ... [详细]
  • 本文详细介绍了如何在PHP中使用Memcached进行数据缓存,包括服务器连接、数据操作、高级功能等。 ... [详细]
  • 本文探讨了如何选择一个合适的序列化版本ID(serialVersionUID),包括使用生成器还是简单的整数,以及在不同情况下应如何处理序列化版本ID。 ... [详细]
  • 如何使用Maven将依赖插件一并打包进JAR文件
    本文详细介绍了在使用Maven构建项目时,如何将所需的依赖插件一同打包进最终的JAR文件中,以避免手动部署依赖库的麻烦。 ... [详细]
  • Go语言实现文件读取与终端输出
    本文介绍如何使用Go语言编写程序,通过命令行参数指定文件路径,读取文件内容并将其输出到控制台。代码示例中包含了错误处理和资源管理的最佳实践。 ... [详细]
  • 本文详细介绍了在Luat OS中如何实现C与Lua的混合编程,包括在C环境中运行Lua脚本、封装可被Lua调用的C语言库,以及C与Lua之间的数据交互方法。 ... [详细]
  • 本文详细探讨了在Java编程语言中,构造函数、静态代码块和构造代码块的执行顺序。首先明确了静态代码块、构造代码块以及构造函数方法体的执行优先级,随后深入分析了构造函数体执行前的具体步骤,包括父类构造器的调用、非静态变量的初始化等。 ... [详细]
author-avatar
求道金林
这个家伙很懒,什么也没留下!
PHP1.CN | 中国最专业的PHP中文社区 | DevBox开发工具箱 | json解析格式化 |PHP资讯 | PHP教程 | 数据库技术 | 服务器技术 | 前端开发技术 | PHP框架 | 开发工具 | 在线工具
Copyright © 1998 - 2020 PHP1.CN. All Rights Reserved | 京公网安备 11010802041100号 | 京ICP备19059560号-4 | PHP1.CN 第一PHP社区 版权所有