Spring Amqp etiketine sahip kayıtlar gösteriliyor. Tüm kayıtları göster
Spring Amqp etiketine sahip kayıtlar gösteriliyor. Tüm kayıtları göster

8 Kasım 2023 Çarşamba

Spring Amqp RabbitMQ Dead Letter Queue

Örnek
Şöyle yaparız
@Configuration
public class RabbitMQConfig { public static final String QUEUE_NAME = "queue-name"; public static final String DLQ_NAME = "dlq-name"; public static final String DLX_NAME = "dlx-name"; @Bean public DirectExchange directExchange() { return new DirectExchange(DIRECT_EXCHANGE); } @Bean public Queue queue() { return new Queue(QUEUE_NAME, true); // Declare the queue as durable } @Bean public Binding binding() { return BindingBuilder.bind(queue()).to(directExchange()).with(QUEUE_NAME); } // DLQ create @Bean public Queue deadLetterQueue() { return new Queue(DLQ_NAME, true); // Declare the DLQ as durable } @Bean public DirectExchange deadLetterExchange() { return new DirectExchange(DLX_NAME); } @Bean public Binding deadLetterBinding() { return BindingBuilder.bind(deadLetterQueue()) .to(deadLetterExchange()).with(DLQ_NAME); } @Bean public Binding queueToDeadLetterExchangeBinding() { return BindingBuilder.bind(queue()) .to(deadLetterExchange()).with(QUEUE_NAME); } }
Şöyle yaparız. Burada DLQ'ya kodla gönderiliyor.
@RabbitListener(queues = "queue-name")
public void handleMessage(String message) {
  try {
    // Process the incoming message
    if (someCondition) {
      throw new Exception("Simulated exception");
    }
    // Message processing succeeded
  } catch (Exception e) {
    // Handle the exception or log it
    System.err.println("Error processing message: " + e.getMessage());
            
    // Send the message to the DLQ
    rabbitTemplate.send("direct.exchange", "dlq-name", new Message(message.getBytes()));
  }
}


13 Ocak 2021 Çarşamba

Spring Amqp RabbitMQ MessageBuilder Sınıfı

Giriş
Şu satırı dahil ederiz
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageBuilder;
import org.springframework.amqp.core.MessageDeliveryMode;
setDeliveryMode metodu
Örnek - PERSISTENT
Şöyle yaparız
Message sysErrMsg = MessageBuilder.withBody(systemErrorLog.getBytes())
                      .setDeliveryMode(MessageDeliveryMode.PERSISTENT)
                      .build();
Örnek - NON_PERSISTENT
Şöyle yaparız
Message sysErrMsg = MessageBuilder.withBody(systemErrorLog.getBytes())
                      .setDeliveryMode(MessageDeliveryMode.NON_PERSISTENT)
                      .build();
setMessageId metodu
Örnek
Şöyle yaparız.
Message message = MessageBuilder.withBody(messageText.getBytes())
            .setMessageId("123")
            .setContentType("application/json")
            .setHeader("foo", "bar")
            .build();

12 Ocak 2021 Salı

Spring Amqp RabbitMQ MessagePostProcessor Arayüzü

Giriş
Şu satırı dahil ederiz
import org.springframework.amqp.core.MessagePostProcessor;
Açıklama şöyle.
Used in several places in the framework, such as AmqpTemplate#convertAndSend(Object, MessagePostProcessor) where it can be used to add/modify headers or properties after the message conversion has been performed.
Örnek - Header Değeri Atamak
Elimizde şöyle bir kod olsun
public class MyMessagePostProcessor implements MessagePostProcessor {

  private final Integer ttl;

  public MyMessagePostProcessor(final Integer ttl) {
    this.ttl = ttl;
  }

  @Override
  public Message postProcessMessage(final Message message) throws AmqpException {
    message.getMessageProperties().getHeaders().put("expiration", ttl.toString());
    return message;
  }
}
Şöyle yaparız. Birinci parametre exchange ismi, ikinci parametre routingkey, üçüncü parametre mesaj verisi, dördüncü parametre MessagePostProcessor nesnesi
inal String message = "message";
final MessagePostProcessor messagePostProcessor = new MyMessagePostProcessor(10000);
rabbitTemplate.convertAndSend("my.queue", "routingKey", message, messagePostProcessor);
Örnek - Property Değeri Atamak - Persistent Olmayan Mesaj
Açıklaması şöyle
The default delivery mode (in MessageProperties) is PERSISTENT.
To make it non-persistent you need to use a convertAndSend(...) method with a MessagePostProcessor to set the deliveryMode property.
Örnek - Property Değeri Atamak
Şöyle yaparızBirinci parametre exchange ismi, ikinci parametre routingkey, üçüncü parametre mesaj verisi, dördüncü parametre MessagePostProcessor nesnesi
import org.springframework.amqp.support.converter.MessageConverter;
import org.springframework.amqp.support.converter.SimpleMessageConverter;


// bad json
template.setMessageConverter(new SimpleMessageConverter());
template.convertAndSend("", "my.queue", "\"routeId\":\"7\"}", m -> {
  m.getMessageProperties().setContentType("application/json");
  return m;
});

27 Mayıs 2020 Çarşamba

Spring Amqp QueueBuilder Sınıfı

Giriş
Şu satırı dahil ederiz
import org.springframework.amqp.core.QueuBuilder;
RabbitMQ yazısına bakabilirsiniz.

Kullanım
- durable() veya nonDurable() metodu ile kuyruk ismi belirtilir.
- Kuyruk ayarları withArgument() metodu ile tek tek belirtilir, veya withArguments() metodu ile topluca belirtilir.
- build() metodu çağrılarak kuyruk yaratılır.
- Daha sonra yaratılan kuyruk AmqpAdmin veya RabbitAdmin sınıfının declaredQueue() metoduna parametre olarak geçilir.

build metodu
Şöyle yaparız.
@Bean
public Queue demoDeadLetterQueue(RabbitAdmin rabbitAdmin) {
  Queue queue = QueueBuilder.durable("demo.queue.dead-letter")
    .build();
    rabbitAdmin.declareQueue(queue)
    return queue;
}
deadLetterExchange metodu
Bir kuyruk ile Dead Letter Exchange bağlanır. Exchange mantığı sadece RabbitMQ'da vardır.

Örnek
Şöyle yaparız.  Bu yöntem yerine withArguments() çağrıları ile de aynı şey yapılabilir.
@Configuration
public class RMQConfiguration implements Serializable {

  String DEAD_LETTER_EXCHANGE_NAME;

  String INCOMING_QUEUE_NAME;

  @Bean
  Queue queue() {
    return QueueBuilder.durable(INCOMING_QUEUE_NAME)
      .deadLetterExchange(DEAD_LETTER_EXCHANGE_NAME)
      .build();
  }
}
Örnek
deadLetterExchange() aslında altta bazı string sabitleri kullanıyor. Aynı şeyi istersek şöyle de yaparız
@Bean
Queue delayQueuePerMessageTTL() {
  return QueueBuilder.durable(DELAY_QUEUE_PER_MESSAGE_TTL_NAME)
    .withArgument("x-dead-letter-exchange", DELAY_EXCHANGE_NAME)
    .withArgument("x-dead-letter-routing-key", DELAY_PROCESS_QUEUE_NAME)
    .build();
}
DLX exchange + queue ikilisini yaratmak için şöyle yaparız
@Bean
Queue delayProcessQueue() {
  return QueueBuilder.durable(DELAY_PROCESS_QUEUE_NAME)
                       .build();
}

@Bean
DirectExchange delayExchange() {
  return new DirectExchange(DELAY_EXCHANGE_NAME);
}

@Bean
Binding dlxBinding(Queue delayProcessQueue, DirectExchange delayExchange) {
  return BindingBuilder.bind(delayProcessQueue)
    .to(delayExchange)
    .with(DELAY_PROCESS_QUEUE_NAME);
}
durable metodu
Parametre olarak kuyruk ismini alır.
Örnek
Şöyle yaparız.
@Bean
public Queue demoQueue(RabbitAdmin rabbitAdmin) {
  Queue queue = QueueBuilder.durable("demo.queue")
    .withArgument("x-dead-letter-exchange", StringUtils.EMPTY)
    .withArgument("x-dead-letter-routing-key", "demo.queue.dead-letter")
    .build();
    queue.setAdminsThatShouldDeclare(rabbitAdmin);
    return queue;
}
Örnek
Şöyle yaparız
@Bean
public Queue createDelayedMessageQueue(RabbitAdmin rabbitAdmin) {
  Queue delayedMessageQueue = QueueBuilder.durable("...")
    .withArgument("x-dead-letter-exchange", "")
    .withArgument("x-dead-letter-routing-key", "...")
    .withArgument("x-queue-mode", "lazy")
    .build();
    delayedMessageQueue.setAdminsThatShouldDeclare(rabbitAdmin);
    return delayedMessageQueue;
}
nonDurable metodu
Parametre olarak kuyruk ismini alır.
Örnek
Şöyle yaparız.
@Bean
public Queue outgoingQueue() {
  Map<String, Object> args = new HashMap<>();
  // The default exchange
  args.put("x-dead-letter-exchange", "");
  // Route to the dead letter queue when the TTL occurs
  args.put("x-dead-letter-routing-key", Constants.DEAD_LETTER_QUEUE);
  // TTL 10 seconds
  args.put("x-message-ttl", 10000);
  args.put("x-max-priority", 10);

  Queue queue =   QueueBuilder.nonDurable(Constants.LIVE_QUEUE)
    .withArguments(args).build();

  amqpAdmin.declareQueue(queue);


  RabbitMqUtil.createShovel(props, appProperties,
     queue.getName(), Constants.MASTER_LIVE_QUEUE);

  return queue;
}
Diğer
"x-dead-letter-exchange" : Dead letter kuyruğuna yazmak için kullanılacak exchange ismi. 
 -Default exchange ise boş string verilir, "x-dead-letter-routing-key" kuyruk ismine sahiptir. Açıklaması şöyle
.. we are using the default exchange (no-name) as the dead letter exchange and using the dead letter queue name as the new routing key. This will work since any queue is bound to the default exchange with the binding key equal to the queue name.
- Direct exchange kullanıyorsak, "x-dead-letter-exchange" kuyruk ismini alır, "x-dead-letter-routing-key" boş string verilir.

"x-dead-letter-routing-key" : Yukarıdaki açıklamaya bakınız

"x-message-ttl" : Mesajın TTL değeri. TTL süresi içinde işlenmezse dead-letter kuyruğuna gönderilir.

"x-max-priority" : Kuyruğa yazılabilecek mesajların alabileceği en yüksek öncelik değeri

"x-queue-mode" : lazy ise mesajlar bellekte tutulmaz, diske yazılır.



Spring Amqp Message Sınıfı

getMessageProperties metodu
Şöyle yaparız
message.getMessageProperties().getXDeathHeader();

21 Mayıs 2020 Perşembe

Spring Amqp RabbitMQ RabbitTemplate Sınıfı

Giriş
Şu satırı dahil ederiz.
import org.springframework.amqp.rabbit.core.RabbitTemplate;
AmqpTemplate sınıfından kalıtır.

constructor
Şöyle yaparız.
ConfigurableApplicationContext context = ...;
RabbitTemplate template = context.getBean(RabbitTemplate.class);
convertAndSend metodu
Örnek - routingKey + Message
Burada default exchange'e gönderilir. Şöyle yaparız.
template.convertAndSend("my.queue", new Foo());
Örnek - exchange + routingKey + message
Şöyle yaparız.
public void send(Blog blog) {
  rabbitTemplate.convertAndSend(config.getExchange(), config.getRoutingkey(), blog);
}
Örnek - exchange + routingKey + message + messagePostProcessor
MessagePostProcessor yazısına taşıdım

send metodu
Örnek
Şöyle yaparız
import org.springframework.amqp.core.*; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @Service public class MessageSenderService { private final RabbitTemplate rabbitTemplate; @Autowired public MessageSenderService(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } public void sendMessage(String message) { MessageProperties properties = new MessageProperties(); // Set message as persistent properties.setDeliveryMode(MessageDeliveryMode.PERSISTENT); Message rabbitMessage = new Message(message.getBytes(), properties); rabbitTemplate.send("my-exchange", "my-routing-key", rabbitMessage); } }
Açıklaması şöyle
By setting both the message delivery mode to PERSISTENT and declaring queues as durable, you ensure that messages sent to RabbitMQ are persisted, and even if RabbitMQ server restarts, the messages will not be lost.
Açıklaması şöyle
Ensure that your RabbitMQ queues are also declared as durable. You can do this by setting the durable attribute to true when declaring queues in your RabbitMQ configuration.

RabbitMQ stores persistent messages in its storage mechanism, which typically consists of regular filesystem storage on the server where RabbitMQ is running. RabbitMQ stores persistent messages on the server’s disk in a location specified by its configuration. This location is usually within the RabbitMQ data directory.
setExchange metodu
RabbitTemplate sınıfının sadece send() metodunu kullanabilmek için önce RabiitTemplat nesnesini sabit olarak bir exchange'e bağlamak gerekir.
Örnek
Şöyle yaparız
@Bean
public RabbitTemplate rubeExchangeTemplate() {
  RabbitTemplate r = new RabbitTemplate(this.rabbitConnectionFactory);
  r.setExchange("rmq-exchange");
  r.setMessageConverter(jsonMessageConverter());
  return r;
}
Örnek
Şöyle yaparızz
@Autowired
private RabbitTemplate   rabbit;

@Autowired
private MessageConverter jsonMessageConverter;

public void produce() {

  rabbit.setExchange("My.Exchange");
  rabbit.setRoutingKey("R.K");
  rabbit.setMessageConverter(jsonMessageConverter);
  MessageProperties props = new MessageProperties();
  props.setExpiration(Long.toString(expiration));
  Message toSend = new Message(message.toString().getBytes(), props);
  rabbit.send(toSend);
}
setMessageConverter metodu
Eğer bir MessageConverter atanmadıysa, RabbitTemplate Java Binary Serialization yöntemini kullanır.
JSON olarak göndermek için Jackson2JsonMessageConverter kullanılır.