Vagrant & ActiveMQ Classic Clusterı
Zaman zaman troubleshootlarda yaşadığımız sorunları doğrudan canlıda tekrarlayamıyoruz ya da 2-3 ayda bir kere denk gelen şeylerin troubleshootunu yapmak güç oluyor. Yine buna benzer bir durum için activemq cluster’ıyla ilgili deneme yanılmalar yapmam gerekti ve kırıp döküp bozacağım bir ortama ihtiyacım oldu. Böyle bir ortamı oluşturmak için de localde Vagrant ile 3 node’lu bir cluster kurmak istedim.
Ön gereksinimler
Vagrant ve Virtualbox kurulu olması gerekiyor. Şu an kullandığım makine windows olduğu için aşağıdaki gibi kurdum. Kurulum adımı sonrası WSL’den devam edeceğim.
winget install HashiCorp.Vagrant
winget install Oracle.VirtualBox
vagrant --version # Vagrant 2.4.9
VBoxManage --version # 7.0.x
Dizi yapısı
WORKDIR=$HOMEDIR/activemq-cluster
tree $WORKDIR :
.
├── files
│ └── apache-activemq-5.16.7
│ └── activemq-node1.xml
│ └── activemq-node2.xml
│ └── activemq-node3.xml
│ └── jetty.xml
├── java
│ └── build.sh
│ └── Consumer.java
│ └── Producer.java
├── provision
│ └── activemq.xml.tpl
│ └── install.sh
└── test
│ └── test.sh
├── README.md
└── Vagrantfile
Vagrant File içeriği
BROKERS = [
{ name: "broker1", ip: "192.168.56.11" },
{ name: "broker2", ip: "192.168.56.12" },
{ name: "broker3", ip: "192.168.56.13" },
]
Vagrant.configure("2") do |config|
# base image
config.vm.box = "ubuntu/jammy64"
config.vm.box_check_update = false
BROKERS.each do |broker|
config.vm.define broker[:name] do |node|
node.vm.hostname = broker[:name]
node.vm.network "private_network", ip: broker[:ip]
# resource definitions
node.vm.provider "virtualbox" do |vb|
vb.name = broker[:name]
vb.memory = 1024
vb.cpus = 1
end
# Binary
node.vm.provision "file",
source: "files/apache-activemq-5.16.7-bin.tar.gz",
destination: "/tmp/apache-activemq-5.16.7-bin.tar.gz"
# copy the configs to nodes
node.vm.provision "file",
source: "files/activemq-node1.xml",
destination: "/tmp/activemq-node1.xml"
node.vm.provision "file",
source: "files/activemq-node2.xml",
destination: "/tmp/activemq-node2.xml"
node.vm.provision "file",
source: "files/activemq-node3.xml",
destination: "/tmp/activemq-node3.xml"
node.vm.provision "file",
source: "files/jetty.xml",
destination: "/tmp/jetty.xml"
# installation script
node.vm.provision "shell",
path: "provision/install.sh",
env: {
"BROKER_NAME" => broker[:name],
"BROKER_IP" => broker[:ip],
}
end
end
end
Node’ların Başlatılması
“$WORKDIR” altındaki vagrantfile’ımızun bulunduğu dizinde vagrant up ile veriyoruz startı. Örnek bir outputun kısaltılmış halini aşağıda bulabilirsiniz, tüm outputu da repomdaki test/ folderı altında mevcut.
vagrant up
==> broker1: Importing base box 'ubuntu/jammy64'...
==> broker1: Matching MAC address for NAT networking...
==> broker1: Setting the name of the VM: broker1
==> broker1: Clearing any previously set network interfaces...
==> broker1: Preparing network interfaces based on configuration...
broker1: Adapter 1: nat
broker1: Adapter 2: hostonly
==> broker1: Forwarding ports...
broker1: 22 (guest) => 2222 (host) (adapter 1)
==> broker1: Running 'pre-boot' VM customizations...
==> broker1: Booting VM...
==> broker1: Waiting for machine to boot. This may take a few minutes...
broker1: SSH address: 127.0.0.1:2222
broker1: SSH username: vagrant
broker1: SSH auth method: private key
broker1: Warning: Connection reset. Retrying...
broker1:
broker1: Vagrant insecure key detected. Vagrant will automatically replace
broker1: this with a newly generated keypair for better security.
broker1:
broker1: Inserting generated public key within guest...
broker1: Removing insecure key from the guest if it's present...
broker1: Key inserted! Disconnecting and reconnecting using new SSH key...
==> broker1: Machine booted and ready!
..
broker1: Guest Additions Version: 6.0.0 r127566
broker1: VirtualBox Version: 7.2
==> broker1: Setting hostname...
==> broker1: Configuring and enabling network interfaces...
==> broker1: Mounting shared folders...
broker1: C:/Ake/workspace/playgrounds/activemq-cluster => /vagrant
==> broker1: Running provisioner: file...
broker1: files/apache-activemq-5.16.7-bin.tar.gz => /tmp/apache-activemq-5.16.7-bin.tar.gz
==> broker1: Running provisioner: file...
broker1: files/activemq-node1.xml => /tmp/activemq-node1.xml
==> broker1: Running provisioner: file...
broker1: files/activemq-node2.xml => /tmp/activemq-node2.xml
==> broker1: Running provisioner: file...
broker1: files/activemq-node3.xml => /tmp/activemq-node3.xml
==> broker1: Running provisioner: file...
broker1: files/jetty.xml => /tmp/jetty.xml
==> broker1: Running provisioner: shell...
broker1: Running: C:/Users/PC/AppData/Local/Temp/vagrant-shell20260803-5760-tt26ug.sh
broker1: >>> [broker1] Installing...
....
broker1: done.
broker1:
broker1: Running kernel seems to be up-to-date.
broker1:
broker1: No services need to be restarted.
broker1:
broker1: No containers need to be restarted.
broker1:
broker1: No user sessions are running outdated binaries.
broker1:
broker1: No VM guests are running outdated hypervisor (qemu) binaries on this host.
broker1: >>> Binary extracting
broker1: Created symlink /etc/systemd/system/multi-user.target.wants/activemq.service → /etc/systemd/system/activemq.service.
broker1: >>> [broker1] Activemq is working!
broker1: Web Konsol : http://192.168.56.11:8161/admin (admin/admin)
broker1: OpenWire : tcp://192.168.56.11:8616
...
VirtualBox’ta instance’ların durumları.

Console’dan kontrol:

VM’e ssh için aşağıdaki komutu kullanabiliriz. Bu komutları yine $WORKDIR altında çalıştırıyorum.
vagrant ssh broker1
vagrant@broker1:/opt/servers$ ls -rltah /opt/servers/activemq/
total 40K
...
drwxr-xr-x 2 activemq activemq 4.0K Aug 3 06:59 conf
drwxr-xr-x 4 activemq activemq 4.0K Aug 3 06:59 data
drwxr-xr-x 4 activemq activemq 4.0K Aug 3 06:59 tmp
Activemq dosyalarımı VM dışında tutuyorum ve onlarda değişiklikler yaptım. Dosyaların güncellenmesi için reload provision yapıyorum.
vagrant reload --provision
Test Cases
Creating Consumer & Producer
Producer
import javax.jms.*;
import org.apache.activemq.ActiveMQConnectionFactory;
/**
* ActiveMQ Producer - Configurable message sending rate
*
*/
public class Producer {
public static void main(String[] args) {
if (args.length < 4) {
System.out.println("Usage: java Producer <brokerUrl> <queueName> <messageCount> <delayMs>");
System.out.println();
System.out.println(" brokerUrl : ActiveMQ broker URL (e.g. tcp://localhost:8616)");
System.out.println(" queueName : Queue name (e.g. TEST.QUEUE)");
System.out.println(" messageCount : Number of messages to send (0 = unlimited)");
System.out.println(" delayMs : Delay between messages in milliseconds");
System.out.println();
System.out.println("Example: java Producer tcp://localhost:8616 TEST.QUEUE 100 1000");
return;
}
String brokerUrl = args[0];
String queueName = args[1];
int messageCount = Integer.parseInt(args[2]);
long delayMs = Long.parseLong(args[3]);
boolean unlimited = (messageCount == 0);
Connection connection = null;
try {
// Create connection
ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory(brokerUrl);
connection = factory.createConnection();
connection.start();
// Create session and producer
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
Destination destination = session.createQueue(queueName);
MessageProducer producer = session.createProducer(destination);
producer.setDeliveryMode(DeliveryMode.PERSISTENT);
System.out.println("=== ActiveMQ Producer ===");
System.out.println("Broker : " + brokerUrl);
System.out.println("Queue : " + queueName);
System.out.println("Messages : " + (unlimited ? "unlimited (Ctrl+C to stop)" : messageCount));
System.out.println("Delay : " + delayMs + " ms");
System.out.println("=========================");
System.out.println();
int sent = 0;
while (unlimited || sent < messageCount) {
sent++;
String text = "Message-" + sent + " | timestamp=" + System.currentTimeMillis();
TextMessage message = session.createTextMessage(text);
message.setIntProperty("messageNumber", sent);
producer.send(message);
System.out.println("[SENT] #" + sent + " -> " + text);
if (delayMs > 0 && (unlimited || sent < messageCount)) {
Thread.sleep(delayMs);
}
}
System.out.println();
System.out.println("Done! Total sent: " + sent);
} catch (Exception e) {
System.err.println("ERROR: " + e.getMessage());
e.printStackTrace();
} finally {
try {
if (connection != null) connection.close();
} catch (JMSException e) {
// ignore
}
}
}
}
Consumer
import javax.jms.*;
import org.apache.activemq.ActiveMQConnectionFactory;
/**
* ActiveMQ Consumer - Configurable message consuming rate
*
*/
public class Consumer {
public static void main(String[] args) {
if (args.length < 4) {
System.out.println("Usage: java Consumer <brokerUrl> <queueName> <messageCount> <delayMs> [prefetch]");
System.out.println();
System.out.println(" brokerUrl : ActiveMQ broker URL (e.g. tcp://localhost:8616)");
System.out.println(" queueName : Queue name (e.g. TEST.QUEUE)");
System.out.println(" messageCount : Number of messages to consume (0 = unlimited)");
System.out.println(" delayMs : Processing delay per message in milliseconds");
System.out.println(" prefetch : (Optional) Prefetch size, default=1 for slow consumers");
System.out.println();
System.out.println("Example: java Consumer tcp://localhost:8616 TEST.QUEUE 50 3000");
System.out.println(" (consume 50 messages, 3 seconds between each)");
return;
}
String brokerUrl = args[0];
String queueName = args[1];
int messageCount = Integer.parseInt(args[2]);
long delayMs = Long.parseLong(args[3]);
int prefetch = args.length >= 5 ? Integer.parseInt(args[4]) : 1;
boolean unlimited = (messageCount == 0);
// Add prefetch to broker URL
String connectionUrl = brokerUrl + "?jms.prefetchPolicy.queuePrefetch=" + prefetch;
Connection connection = null;
try {
// Create connection
ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory(connectionUrl);
connection = factory.createConnection();
connection.start();
// Create session and consumer
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
Destination destination = session.createQueue(queueName);
MessageConsumer consumer = session.createConsumer(destination);
System.out.println("=== ActiveMQ Consumer ===");
System.out.println("Broker : " + brokerUrl);
System.out.println("Queue : " + queueName);
System.out.println("Messages : " + (unlimited ? "unlimited (Ctrl+C to stop)" : messageCount));
System.out.println("Delay : " + delayMs + " ms (processing time per message)");
System.out.println("Prefetch : " + prefetch);
System.out.println("=========================");
System.out.println();
System.out.println("Waiting for messages...");
System.out.println();
int consumed = 0;
while (unlimited || consumed < messageCount) {
// Wait up to 10 seconds for a message
Message message = consumer.receive(10000);
if (message == null) {
System.out.println("[TIMEOUT] No message received in 10s, still waiting...");
continue;
}
consumed++;
String body = "";
if (message instanceof TextMessage) {
body = ((TextMessage) message).getText();
} else {
body = message.toString();
}
System.out.println("[RECV] #" + consumed + " <- " + body);
// Simulate slow processing
if (delayMs > 0) {
System.out.println(" Processing... (" + delayMs + " ms)");
Thread.sleep(delayMs);
System.out.println(" Done.");
}
}
System.out.println();
System.out.println("Done! Total consumed: " + consumed);
} catch (Exception e) {
System.err.println("ERROR: " + e.getMessage());
e.printStackTrace();
} finally {
try {
if (connection != null) connection.close();
} catch (JMSException e) {
// ignore
}
}
}
}
wsl kullanmanın avantajıyla .java dosyalarımızı executable olması için build ediyoruz. Bize Consumer.class ve Producer.class dosyalarını oluşturuyor. Ayrıca activemq librarysinin kullanılması için de classpath’i point etmemiz gerekiyor. Bir node’a mesaj gönderip diğerinden okuyorum. Burada farklı farklı caseler denemem gerekmişti. Buradan sonrası istediğimiz senaryoyu oluşturmaya kalıyor.
cd "$WORKDIR/java"
./build.sh
export JAVA_HOME="/mnt/c/Ake/Apps/java/current"
export CLASSPATH="/mnt/c/Ake/workspace/playgrounds/activemq-cluster/files/apache-activemq-5.16.7/activemq-all-5.16.7.jar":."
test1
Produce message
$JAVA_HOME/bin/java Producer tcp://192.168.56.12:8616 TEST.QUEUE 1000 1000
=== ActiveMQ Producer ===
Broker : tcp://192.168.56.12:8616
Queue : TEST.QUEUE
Messages : 1000
Delay : 1000 ms
=========================
[SENT] #1 -> Message-1 | timestamp=1785782311394
[SENT] #2 -> Message-2 | timestamp=1785782312436
[SENT] #3 -> Message-3 | timestamp=1785782313443
Consume message
$JAVA_HOME/bin/java Consumer tcp://192.168.56.11:8616 TEST.QUEUE 100 50
=== ActiveMQ Consumer ===
Broker : tcp://192.168.56.11:8616
Queue : TEST.QUEUE
Messages : 100
Delay : 50 ms (processing time per message)
Prefetch : 1
=========================
Waiting for messages...
[RECV] #1 <- Message-1 | timestamp=1785781846885
Processing... (50 ms)
Done.
[RECV] #2 <- Message-2 | timestamp=1785781847913
Processing... (50 ms)
Done.

Tüm dosyalara repomdan ulaşabilirsiniz. github
İyi günler dilerim.