< ahmetkaygisiz />

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ı. VirtualBox-Status

Console’dan kontrol: Activemq Console

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.

Queue Status

Tüm dosyalara repomdan ulaşabilirsiniz. github

İyi günler dilerim.