ActiveMQ 訊息持久化到資料庫
阿新 • • 發佈:2019-01-06
1:前言
這一段給公司開發訊息匯流排有機會研究ActiveMQ,今天撰文給大家介紹一下他的持久化訊息。本文只介紹三種方式,分別是持久化為檔案,MYSql,Oracle。下面逐一介紹。
A:持久化為檔案
這個你裝ActiveMQ時預設就是這種,只要你設定訊息為持久化就可以了。涉及到的配置和程式碼有
<persistenceAdapter><kahaDB directory="${activemq.base}/data/kahadb"/></persistenceAdapter>producer.Send(request, MsgDeliveryMode.Persistent, level, TimeSpan.MinValue);
B:持久化為MySql
你首先需要把MySql的驅動放到ActiveMQ的Lib目錄下,我用的檔名字是:mysql-connector-java-5.0.4-bin.jar
接下來你修改配置檔案
<persistenceAdapter>
<jdbcPersistenceAdapter dataDirectory="${activemq.base}/data" dataSource="#derby-ds"/>
</persistenceAdapter>
在配置檔案中的broker節點外增加
<bean id="derby-ds"class="org.apache.commons.dbcp.BasicDataSource" destroy-method="close"> <property name="driverClassName" value="com.mysql.jdbc.Driver"/> <property name="url" value="jdbc:mysql://localhost/activemq?relaxAutoCommit=true"/> <property name="username" value="activemq"/> <property name="password" value="activemq"/> <property name="maxActive" value="200"/> <property name="poolPreparedStatements" value="true"/> </bean>
從配置中可以看出資料庫的名稱是activemq,你需要手動在MySql中增加這個庫。
然後重新啟動訊息佇列,你會發現多了3張表
1:activemq_acks
2:activemq_lock
3:activemq_msgs
C:持久化為Oracle
和持久化為MySql一樣。這裡我說兩點
1;在ActiveMQ安裝資料夾裡的Lib資料夾中增加Oracle的JDBC驅動。驅動檔案位於Oracle客戶端安裝檔案中的product\11.1.0\client_1\jdbc\lib資料夾下。
2:
<bean id="derby-ds"class="org.apache.commons.dbcp.BasicDataSource" destroy-method="close"> <property name="driverClassName" value="oracle.jdbc.driver.OracleDriver"/> <property name="url" value="jdbc:oracle:thin:@10.53.132.47:1521:test"/> <property name="username" value="qdcommu"/> <property name="password" value="qdcommu"/> <property name="maxActive" value="200"/> <property name="poolPreparedStatements" value="true"/> </bean>
這裡的jdbc:oracle:thin:@10.53.132.47:1521:test按照自己實際情況設定一下就可以了,特別注意的是test是SID即服務名稱而不是TNS中配置的節點名。各位同學只需要替換IP,埠和這個SID就可以了。訊息消費者的事先程式碼: 訊息消費者的事先程式碼:
package easyway.activemq.app;
import javax.jms.Connection;
import javax.jms.Destination;
import javax.jms.Message;
import javax.jms.MessageConsumer;
import javax.jms.Session;
import org.apache.activemq.ActiveMQConnectionFactory;
import org.apache.activemq.broker.BrokerService;
import org.apache.log4j.LogManager;
import org.apache.log4j.Logger;
/***
* 訊息持久化到資料庫
* @author longgangbai
*/
public class MessageCustomer {
private static Logger logger=LogManager.getLogger(MessageProductor.class);
private String username=ActiveMQConnectionFactory.DEFAULT_USER;
private String password=ActiveMQConnectionFactory.DEFAULT_PASSWORD;
private String url=ActiveMQConnectionFactory.DEFAULT_BROKER_BIND_URL;
private static String QUEUENAME="ActiveMQ.QUEUE";
protected static final int messagesExpected = 10;
protected ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(
url+"?jms.prefetchPolicy.all=0&jms.redeliveryPolicy.maximumRedeliveries="+messagesExpected);
/***
* 建立Broker服務物件
* @return
* @throws Exception
*/
public BrokerService createBroker()throws Exception{
BrokerService broker=new BrokerService();
broker.addConnector(url);
return broker;
}
/**
* 啟動BrokerService程序
* @throws Exception
*/
public void init() throws Exception{
BrokerService brokerService=createBroker();
brokerService.start();
}
/**
* 接收的資訊
* @return
* @throws Exception
*/
public int receiveMessage() throws Exception{
Connection connection=connectionFactory.createConnection();
connection.start();
Session session = connection.createSession(true, Session.SESSION_TRANSACTED);
return receiveMessages(messagesExpected,session);
}
/**
* 接受資訊的方法
* @param messagesExpected
* @param session
* @return
* @throws Exception
*/
protected int receiveMessages(int messagesExpected, Session session) throws Exception {
int messagesReceived = 0;
for (int i=0; i<messagesExpected; i++) {
Destination destination = session.createQueue(QUEUENAME);
MessageConsumer consumer = session.createConsumer(destination);
Message message = null;
try {
logger.debug("Receiving message " + (messagesReceived+1) + " of " + messagesExpected);
message = consumer.receive(2000);
logger.info("Received : " + message);
if (message != null) {
session.commit();
messagesReceived++;
}
} catch (Exception e) {
logger.debug("Caught exception " + e);
session.rollback();
} finally {
if (consumer != null) {
consumer.close();
}
}
}
return messagesReceived;
}
public String getPassword() {
return password;
}
public void setPassword(String password) {
this.password = password;
}
public String getUrl() {
return url;
}
public void setUrl(String url) {
this.url = url;
}
public String getUsername() {
return username;
}
public void setUsername(String username) {
this.username = username;
}
}
訊息生產者的程式碼:
package easyway.activemq.app;
import java.io.File;
import java.util.Properties;
import javax.jms.Connection;
import javax.jms.DeliveryMode;
import javax.jms.Destination;
import javax.jms.JMSException;
import javax.jms.MessageProducer;
import javax.jms.Session;
import javax.sql.DataSource;
import org.apache.activemq.ActiveMQConnectionFactory;
import org.apache.activemq.broker.BrokerService;
import org.apache.activemq.store.jdbc.adapter.MySqlJDBCAdapter;
import org.apache.commons.dbcp.BasicDataSourceFactory;
import org.apache.log4j.LogManager;
import org.apache.log4j.Logger;
import easyway.activemq.app.utils.BrokenPersistenceAdapter;
/**
* 訊息持久化到資料庫
* @author longgangbai
*
*/
public class MessageProductor {
private static Logger logger=LogManager.getLogger(MessageProductor.class);
private String username=ActiveMQConnectionFactory.DEFAULT_USER;
private String password=ActiveMQConnectionFactory.DEFAULT_PASSWORD;
private String url=ActiveMQConnectionFactory.DEFAULT_BROKER_BIND_URL;
private static String queueName="ActiveMQ.QUEUE";
private BrokerService brokerService;
protected static final int messagesExpected = 10;
protected ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(
"tcp://localhost:61617?jms.prefetchPolicy.all=0&jms.redeliveryPolicy.maximumRedeliveries="+messagesExpected);
/***
* 建立Broker服務物件
* @return
* @throws Exception
*/
public BrokerService createBroker()throws Exception{
BrokerService broker=new BrokerService();
BrokenPersistenceAdapter jdbc=createBrokenPersistenceAdapter();
broker.setPersistenceAdapter(jdbc);
jdbc.setDataDirectory(System.getProperty("user.dir")+File.separator+"data"+File.separator);
jdbc.setAdapter(new MySqlJDBCAdapter());
broker.setPersistent(true);
broker.addConnector("tcp://localhost:61617");
//broker.addConnector(ActiveMQConnectionFactory.DEFAULT_BROKER_BIND_URL);
return broker;
}
/**
* 建立Broken的持久化介面卡
* @return
* @throws Exception
*/
public BrokenPersistenceAdapter createBrokenPersistenceAdapter() throws Exception{
BrokenPersistenceAdapter jdbc=new BrokenPersistenceAdapter();
DataSource datasource=createDataSource();
jdbc.setDataSource(datasource);
jdbc.setUseDatabaseLock(false);
//jdbc.deleteAllMessages();
return jdbc;
}
/**
* 建立資料來源
* @return
* @throws Exception
*/
public DataSource createDataSource() throws Exception{
Properties props=new Properties();
props.put("driverClassName", "com.mysql.jdbc.Driver");
props.put("url", "jdbc:mysql://localhost:3306/activemq");
props.put("username", "root");
props.put("password", "root");
DataSource datasource=BasicDataSourceFactory.createDataSource(props);
return datasource;
}
/**
* 啟動BrokerService程序
* @throws Exception
*/
public void init() throws Exception{
createBrokerService();
brokerService.start();
}
public BrokerService createBrokerService() throws Exception{
if(brokerService==null){
brokerService=createBroker();
}
return brokerService;
}
public void sendMessage() throws JMSException{
Connection connection=connectionFactory.createConnection();
connection.start();
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
Destination destination = session.createQueue(queueName);
MessageProducer producer = session.createProducer(destination);
producer.setDeliveryMode(DeliveryMode.PERSISTENT);
for(int i=0;i<messagesExpected;i++){
logger.debug("Sending message " + (i+1) + " of " + messagesExpected);
producer.send(session.createTextMessage("test message " + (i+1)));
}
connection.close();
}
public String getPassword() {
return password;
}
public void setPassword(String password) {
this.password = password;
}
public String getUrl() {
return url;
}
public void setUrl(String url) {
this.url = url;
}
public String getUsername() {
return username;
}
public void setUsername(String username) {
this.username = username;
}
}
持久化介面卡類
package easyway.activemq.app.utils;
import java.io.IOException;
import org.apache.activemq.broker.ConnectionContext;
import org.apache.activemq.store.jdbc.JDBCPersistenceAdapter;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class BrokenPersistenceAdapter extends JDBCPersistenceAdapter {
private final Logger LOG = LoggerFactory.getLogger(BrokenPersistenceAdapter.class);
private boolean shouldBreak = false;
@Override
public void commitTransaction(ConnectionContext context) throws IOException {
if ( shouldBreak ) {
LOG.warn("Throwing exception on purpose");
throw new IOException("Breaking on purpose");
}
LOG.debug("in commitTransaction");
super.commitTransaction(context);
}
public void setShouldBreak(boolean shouldBreak) {
this.shouldBreak = shouldBreak;
}
}
測測試程式碼如下:
package easyway.activemq.app.test;
import easyway.activemq.app.MessageProductor;
public class MessageProductorTest {
public static void main(String[] args) throws Exception {
MessageProductor productor =new MessageProductor();
productor.init();
productor.sendMessage();
//productor.createBrokerService().stop();
}
}
package easyway.activemq.app.test;
import easyway.activemq.app.MessageCustomer;
public class MessageCustomerTest {
public static void main(String[] args) throws Exception {
MessageCustomer customer=new MessageCustomer();
//customer.init(); //當兩臺機器在不同的伺服器上啟動客戶端的broker程序
customer.receiveMessage();
}
}
備註:執行過程為:首先執行MessageProductorTest,MessageCustomerTest。
mysql資料庫activemq必須存在。關於訊息持久化的表結構如下: