Skip to content

Instantly share code, notes, and snippets.

@chirino
Created May 31, 2013 18:21
Show Gist options
  • Select an option

  • Save chirino/5686895 to your computer and use it in GitHub Desktop.

Select an option

Save chirino/5686895 to your computer and use it in GitHub Desktop.
commit f7f808c6e477822e3b85a9b0ab1d73e6d033966f
Author: Hiram Chirino <hiram@hiramchirino.com>
Date: Fri May 31 14:17:57 2013 -0400
Fix for AMQ-4563: Changes the KahaDB store to index messages by the broker sequence id instead of the message id since that can change depending on how the message was encoded.
diff --git a/activemq-client/src/main/java/org/apache/activemq/command/MessageId.java b/activemq-client/src/main/java/org/apache/activemq/command/MessageId.java
index 9981460..f1c7cbc 100755
--- a/activemq-client/src/main/java/org/apache/activemq/command/MessageId.java
+++ b/activemq-client/src/main/java/org/apache/activemq/command/MessageId.java
@@ -105,6 +105,10 @@ public class MessageId implements DataStructure, Comparable<MessageId> {
return hashCode;
}
+ public String toStringBrokerSequenceId() {
+ return Long.toHexString(brokerSequenceId);
+ }
+
public String toString() {
if (key == null) {
key = producerId.toString() + ":" + producerSequenceId;
diff --git a/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/KahaDBStore.java b/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/KahaDBStore.java
index 5d8bea0..5848566 100644
--- a/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/KahaDBStore.java
+++ b/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/KahaDBStore.java
@@ -423,7 +423,7 @@ public class KahaDBStore extends MessageDatabase implements PersistenceAdapter {
public void addMessage(ConnectionContext context, Message message) throws IOException {
KahaAddMessageCommand command = new KahaAddMessageCommand();
command.setDestination(dest);
- command.setMessageId(message.getMessageId().toString());
+ command.setMessageId(message.getMessageId().toStringBrokerSequenceId());
command.setTransactionInfo(transactionIdTransformer.transform(message.getTransactionId()));
command.setPriority(message.getPriority());
command.setPrioritySupported(isPrioritizedMessages());
@@ -436,7 +436,7 @@ public class KahaDBStore extends MessageDatabase implements PersistenceAdapter {
public void removeMessage(ConnectionContext context, MessageAck ack) throws IOException {
KahaRemoveMessageCommand command = new KahaRemoveMessageCommand();
command.setDestination(dest);
- command.setMessageId(ack.getLastMessageId().toString());
+ command.setMessageId(ack.getLastMessageId().toStringBrokerSequenceId());
command.setTransactionInfo(transactionIdTransformer.transform(ack.getTransactionId()));
org.apache.activemq.util.ByteSequence packet = wireFormat.marshal(ack);
@@ -451,7 +451,7 @@ public class KahaDBStore extends MessageDatabase implements PersistenceAdapter {
}
public Message getMessage(MessageId identity) throws IOException {
- final String key = identity.toString();
+ final String key = identity.toStringBrokerSequenceId();
// Hopefully one day the page file supports concurrent read
// operations... but for now we must
@@ -590,7 +590,7 @@ public class KahaDBStore extends MessageDatabase implements PersistenceAdapter {
@Override
public void setBatch(MessageId identity) throws IOException {
try {
- final String key = identity.toString();
+ final String key = identity.toStringBrokerSequenceId();
lockAsyncJobQueue();
// Hopefully one day the page file supports concurrent read
@@ -707,7 +707,7 @@ public class KahaDBStore extends MessageDatabase implements PersistenceAdapter {
KahaRemoveMessageCommand command = new KahaRemoveMessageCommand();
command.setDestination(dest);
command.setSubscriptionKey(subscriptionKey);
- command.setMessageId(messageId.toString());
+ command.setMessageId(messageId.toStringBrokerSequenceId());
command.setTransactionInfo(ack != null ? transactionIdTransformer.transform(ack.getTransactionId()) : null);
if (ack != null && ack.isUnmatchedAck()) {
command.setAck(UNMATCHED);
diff --git a/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/MessageDatabase.java b/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/MessageDatabase.java
index 1a1c718..4373e26 100644
--- a/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/MessageDatabase.java
+++ b/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/MessageDatabase.java
@@ -2268,7 +2268,7 @@ public abstract class MessageDatabase extends ServiceSupport implements BrokerSe
this.indexLock.writeLock().lock();
try {
for (MessageAck ack : acks) {
- ackedAndPrepared.add(ack.getLastMessageId().toString());
+ ackedAndPrepared.add(ack.getLastMessageId().toStringBrokerSequenceId());
}
} finally {
this.indexLock.writeLock().unlock();
@@ -2280,7 +2280,7 @@ public abstract class MessageDatabase extends ServiceSupport implements BrokerSe
this.indexLock.writeLock().lock();
try {
for (MessageAck ack : acks) {
- ackedAndPrepared.remove(ack.getLastMessageId().toString());
+ ackedAndPrepared.remove(ack.getLastMessageId().toStringBrokerSequenceId());
}
} finally {
this.indexLock.writeLock().unlock();
diff --git a/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/TempKahaDBStore.java b/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/TempKahaDBStore.java
index 3320f89..d903ce1 100644
--- a/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/TempKahaDBStore.java
+++ b/activemq-kahadb-store/src/main/java/org/apache/activemq/store/kahadb/TempKahaDBStore.java
@@ -135,14 +135,14 @@ public class TempKahaDBStore extends TempMessageDatabase implements PersistenceA
public void addMessage(ConnectionContext context, Message message) throws IOException {
KahaAddMessageCommand command = new KahaAddMessageCommand();
command.setDestination(dest);
- command.setMessageId(message.getMessageId().toString());
+ command.setMessageId(message.getMessageId().toStringBrokerSequenceId());
processAdd(command, message.getTransactionId(), wireFormat.marshal(message));
}
public void removeMessage(ConnectionContext context, MessageAck ack) throws IOException {
KahaRemoveMessageCommand command = new KahaRemoveMessageCommand();
command.setDestination(dest);
- command.setMessageId(ack.getLastMessageId().toString());
+ command.setMessageId(ack.getLastMessageId().toStringBrokerSequenceId());
processRemove(command, ack.getTransactionId());
}
@@ -153,7 +153,7 @@ public class TempKahaDBStore extends TempMessageDatabase implements PersistenceA
}
public Message getMessage(MessageId identity) throws IOException {
- final String key = identity.toString();
+ final String key = identity.toStringBrokerSequenceId();
// Hopefully one day the page file supports concurrent read operations... but for now we must
// externally synchronize...
@@ -241,7 +241,7 @@ public class TempKahaDBStore extends TempMessageDatabase implements PersistenceA
@Override
public void setBatch(MessageId identity) throws IOException {
- final String key = identity.toString();
+ final String key = identity.toStringBrokerSequenceId();
// Hopefully one day the page file supports concurrent read operations... but for now we must
// externally synchronize...
@@ -282,7 +282,7 @@ public class TempKahaDBStore extends TempMessageDatabase implements PersistenceA
KahaRemoveMessageCommand command = new KahaRemoveMessageCommand();
command.setDestination(dest);
command.setSubscriptionKey(subscriptionKey(clientId, subscriptionName));
- command.setMessageId(messageId.toString());
+ command.setMessageId(messageId.toStringBrokerSequenceId());
// We are not passed a transaction info.. so we can't participate in a transaction.
// Looks like a design issue with the TopicMessageStore interface. Also we can't recover the original ack
// to pass back to the XA recover method.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment