mirror of https://github.com/apache/activemq.git
merging 885841 - https://issues.apache.org/activemq/browse/AMQ-2519 - duplicate messages and jdbc store
git-svn-id: https://svn.apache.org/repos/asf/activemq/branches/activemq-5.3@885843 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
parent
55dcb0e3e0
commit
9360f50b49
|
@ -18,6 +18,8 @@ package org.apache.activemq.store.jdbc;
|
||||||
|
|
||||||
import java.io.IOException;
|
import java.io.IOException;
|
||||||
import java.sql.SQLException;
|
import java.sql.SQLException;
|
||||||
|
import java.util.HashMap;
|
||||||
|
import java.util.Map;
|
||||||
import java.util.concurrent.atomic.AtomicLong;
|
import java.util.concurrent.atomic.AtomicLong;
|
||||||
|
|
||||||
import org.apache.activemq.broker.ConnectionContext;
|
import org.apache.activemq.broker.ConnectionContext;
|
||||||
|
@ -25,24 +27,28 @@ import org.apache.activemq.command.ActiveMQDestination;
|
||||||
import org.apache.activemq.command.Message;
|
import org.apache.activemq.command.Message;
|
||||||
import org.apache.activemq.command.MessageAck;
|
import org.apache.activemq.command.MessageAck;
|
||||||
import org.apache.activemq.command.MessageId;
|
import org.apache.activemq.command.MessageId;
|
||||||
import org.apache.activemq.store.MessageRecoveryListener;
|
import org.apache.activemq.command.ProducerId;
|
||||||
import org.apache.activemq.store.MessageStore;
|
|
||||||
import org.apache.activemq.store.AbstractMessageStore;
|
import org.apache.activemq.store.AbstractMessageStore;
|
||||||
import org.apache.activemq.usage.MemoryUsage;
|
import org.apache.activemq.store.MessageRecoveryListener;
|
||||||
|
import org.apache.activemq.store.jdbc.adapter.DefaultJDBCAdapter;
|
||||||
import org.apache.activemq.util.ByteSequence;
|
import org.apache.activemq.util.ByteSequence;
|
||||||
import org.apache.activemq.util.ByteSequenceData;
|
import org.apache.activemq.util.ByteSequenceData;
|
||||||
import org.apache.activemq.util.IOExceptionSupport;
|
import org.apache.activemq.util.IOExceptionSupport;
|
||||||
import org.apache.activemq.wireformat.WireFormat;
|
import org.apache.activemq.wireformat.WireFormat;
|
||||||
|
import org.apache.commons.logging.Log;
|
||||||
|
import org.apache.commons.logging.LogFactory;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* @version $Revision: 1.10 $
|
* @version $Revision: 1.10 $
|
||||||
*/
|
*/
|
||||||
public class JDBCMessageStore extends AbstractMessageStore {
|
public class JDBCMessageStore extends AbstractMessageStore {
|
||||||
|
|
||||||
|
private static final Log LOG = LogFactory.getLog(JDBCMessageStore.class);
|
||||||
protected final WireFormat wireFormat;
|
protected final WireFormat wireFormat;
|
||||||
protected final JDBCAdapter adapter;
|
protected final JDBCAdapter adapter;
|
||||||
protected final JDBCPersistenceAdapter persistenceAdapter;
|
protected final JDBCPersistenceAdapter persistenceAdapter;
|
||||||
protected AtomicLong lastMessageId = new AtomicLong(-1);
|
protected AtomicLong lastMessageId = new AtomicLong(-1);
|
||||||
|
protected Map<ProducerId, Long> addedMessages = new HashMap<ProducerId, Long>();
|
||||||
|
|
||||||
public JDBCMessageStore(JDBCPersistenceAdapter persistenceAdapter, JDBCAdapter adapter, WireFormat wireFormat, ActiveMQDestination destination) {
|
public JDBCMessageStore(JDBCPersistenceAdapter persistenceAdapter, JDBCAdapter adapter, WireFormat wireFormat, ActiveMQDestination destination) {
|
||||||
super(destination);
|
super(destination);
|
||||||
|
@ -53,22 +59,32 @@ public class JDBCMessageStore extends AbstractMessageStore {
|
||||||
|
|
||||||
public void addMessage(ConnectionContext context, Message message) throws IOException {
|
public void addMessage(ConnectionContext context, Message message) throws IOException {
|
||||||
|
|
||||||
|
MessageId messageId = message.getMessageId();
|
||||||
|
Long lastAddedMessage = addedMessages.get(messageId.getProducerId());
|
||||||
|
if (lastAddedMessage != null && lastAddedMessage >= messageId.getProducerSequenceId()) {
|
||||||
|
if (LOG.isDebugEnabled()) {
|
||||||
|
LOG.debug("Message " + message + " already added to the database. Skipping.");
|
||||||
|
}
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
// Serialize the Message..
|
// Serialize the Message..
|
||||||
byte data[];
|
byte data[];
|
||||||
try {
|
try {
|
||||||
ByteSequence packet = wireFormat.marshal(message);
|
ByteSequence packet = wireFormat.marshal(message);
|
||||||
data = ByteSequenceData.toByteArray(packet);
|
data = ByteSequenceData.toByteArray(packet);
|
||||||
} catch (IOException e) {
|
} catch (IOException e) {
|
||||||
throw IOExceptionSupport.create("Failed to broker message: " + message.getMessageId() + " in container: " + e, e);
|
throw IOExceptionSupport.create("Failed to broker message: " + messageId + " in container: " + e, e);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get a connection and insert the message into the DB.
|
// Get a connection and insert the message into the DB.
|
||||||
TransactionContext c = persistenceAdapter.getTransactionContext(context);
|
TransactionContext c = persistenceAdapter.getTransactionContext(context);
|
||||||
try {
|
try {
|
||||||
adapter.doAddMessage(c, message.getMessageId(), destination, data, message.getExpiration());
|
adapter.doAddMessage(c, messageId, destination, data, message.getExpiration());
|
||||||
|
addedMessages.put(messageId.getProducerId(), messageId.getProducerSequenceId());
|
||||||
} catch (SQLException e) {
|
} catch (SQLException e) {
|
||||||
JDBCPersistenceAdapter.log("JDBC Failure: ", e);
|
JDBCPersistenceAdapter.log("JDBC Failure: ", e);
|
||||||
throw IOExceptionSupport.create("Failed to broker message: " + message.getMessageId() + " in container: " + e, e);
|
throw IOExceptionSupport.create("Failed to broker message: " + messageId + " in container: " + e, e);
|
||||||
} finally {
|
} finally {
|
||||||
c.close();
|
c.close();
|
||||||
}
|
}
|
||||||
|
|
|
@ -20,6 +20,11 @@ import java.io.IOException;
|
||||||
|
|
||||||
import junit.framework.AssertionFailedError;
|
import junit.framework.AssertionFailedError;
|
||||||
|
|
||||||
|
import org.apache.activemq.broker.ConnectionContext;
|
||||||
|
import org.apache.activemq.command.ActiveMQQueue;
|
||||||
|
import org.apache.activemq.command.ActiveMQTextMessage;
|
||||||
|
import org.apache.activemq.command.MessageId;
|
||||||
|
import org.apache.activemq.store.MessageStore;
|
||||||
import org.apache.activemq.store.PersistenceAdapter;
|
import org.apache.activemq.store.PersistenceAdapter;
|
||||||
import org.apache.activemq.store.PersistenceAdapterTestSupport;
|
import org.apache.activemq.store.PersistenceAdapterTestSupport;
|
||||||
import org.apache.derby.jdbc.EmbeddedDataSource;
|
import org.apache.derby.jdbc.EmbeddedDataSource;
|
||||||
|
@ -42,14 +47,4 @@ public class JDBCPersistenceAdapterTest extends PersistenceAdapterTestSupport {
|
||||||
return jdbc;
|
return jdbc;
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
|
||||||
public void testStoreCanHandleDupMessages() throws Exception {
|
|
||||||
try {
|
|
||||||
super.testStoreCanHandleDupMessages();
|
|
||||||
fail("We expect this test to fail as it would be too expensive to add additional " +
|
|
||||||
"unique constraints in the JDBC implementation to detect the duplicate messages.");
|
|
||||||
} catch (AssertionFailedError expected) {
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
Loading…
Reference in New Issue