Hi guys!
Looks like i found solution for
https://issues.apache.org/jira/projects/JAMES/issues/JAMES-4192
We can use
Artemis _AMQ_SCHED_DELIVERY
If you think i found it right - tell me i can add this to
https://github.com/apache/james-project/pull/3194
or not?
================================================================================
JAMES-4192: Resolving the Delayed Mail Paging Deadlock in Apache ActiveMQ
Artemis
================================================================================
1. PROBLEM ROOT CAUSE ANALYSIS
--------------------------------------------------------------------------------
In Apache James `JMSCacheableMailQueue`, message delay scheduling is currently
implemented using a custom message property `JAMES_NEXT_DELIVERY` combined with
a
JMS Message Selector on the consumer side:
String selector = "JAMES_NEXT_DELIVERY <= " + System.currentTimeMillis()
+ " OR FORCE_DELIVERY = true";
consumer = session.createConsumer(queue, selector);
Why this causes head-of-line blocking and deadlocks in message brokers:
- In ActiveMQ Classic:
The broker scans up to `maxPageSize` (default: 200) messages into memory.
If all 200 scanned messages are delayed retries, the selector rejects them
all.
The broker stops searching deeper into the queue, leaving ready messages stuck
behind delayed ones (head-of-line blocking).
- In ActiveMQ Artemis:
Artemis only evaluates JMS message selectors against messages held in memory,
NOT against messages paged out to disk. Once message paging begins, memory
fills up with delayed messages rejected by `JAMES_NEXT_DELIVERY <= now`.
Because no messages are consumed, no space is freed to page in subsequent
messages from disk. Ready messages residing in page files stay stuck
indefinitely.
2. THE SOLUTION FOR ARTEMIS: NATIVE SCHEDULED DELIVERY (JMS 2.0)
--------------------------------------------------------------------------------
Artemis has a dedicated native scheduled message delivery engine built into its
core
(`ScheduledDeliveryHandler`). Instead of keeping delayed messages in the visible
queue and relying on consumer selectors:
1. Messages are enqueued with a delivery delay via JMS 2.0
`producer.setDeliveryDelay(ms)`
or native Artemis property `_AMQ_SCHED_DELIVERY`.
2. Artemis intercepts these messages and places them in a dedicated priority
scheduler
tree. They are NOT placed in the visible queue or page files, and do NOT
block consumers.
3. When the scheduled timestamp arrives, Artemis automatically moves the
message into
the ready queue.
4. The consumer reads from the queue with NO SELECTOR
(`session.createConsumer(queue)`).
It consumes messages at full wire speed without any paging or evaluation
bottlenecks.
3. CODE IMPLEMENTATION EXAMPLE
--------------------------------------------------------------------------------
Below is an example showing how `ArtemisMailQueue` replaces the selector-based
pattern with native scheduled delivery.
```java
package org.apache.james.queue.artemis;
import java.time.Duration;
import java.util.Map;
import jakarta.jms.Connection;
import jakarta.jms.ConnectionFactory;
import jakarta.jms.JMSException;
import jakarta.jms.Message;
import jakarta.jms.MessageConsumer;
import jakarta.jms.MessageProducer;
import jakarta.jms.Queue;
import jakarta.jms.Session;
import org.apache.james.queue.api.MailQueue;
import org.apache.james.queue.api.MailQueueItem;
import org.apache.james.queue.api.MailQueueName;
import org.apache.james.queue.jms.JMSCacheableMailQueue;
import org.apache.mailet.Mail;
public class ArtemisMailQueue extends JMSCacheableMailQueue {
/**
* Artemis native scheduled delivery header key.
* Value is absolute epoch timestamp in milliseconds when message should be
delivered.
* Using this is equivalent to JMS 2.0 setDeliveryDelay().
*/
public static final String AMQ_SCHEDULED_DELIVERY = "_AMQ_SCHED_DELIVERY";
private final Connection connection;
private final MailQueueName queueName;
public ArtemisMailQueue(ConnectionFactory connectionFactory,
MailQueueName queueName) throws JMSException {
super(connectionFactory, queueName);
this.connection = connectionFactory.createConnection();
this.queueName = queueName;
}
/**
* Enqueue a mail with an optional delay.
*
* Instead of setting JAMES_NEXT_DELIVERY property for consumer-side
selector filtering,
* we configure native Artemis delivery delay.
*/
@Override
public void enQueue(Mail mail, Duration delay) throws MailQueueException {
try (Session session = connection.createSession(false,
Session.AUTO_ACKNOWLEDGE)) {
Queue queue = session.createQueue(queueName.asString());
try (MessageProducer producer = session.createProducer(queue)) {
long delayMillis = (delay != null && !delay.isNegative()) ?
delay.toMillis() : 0L;
if (delayMillis > 0) {
// Option A: Standard JMS 2.0 / Jakarta Messaging API
producer.setDeliveryDelay(delayMillis);
// Option B: Artemis-specific property (useful when
producer pooling or compatibility wrapper is used)
// long deliverAt = System.currentTimeMillis() +
delayMillis;
// message.setLongProperty(AMQ_SCHEDULED_DELIVERY,
deliverAt);
}
Map<String, Object> props = getJMSProperties(mail, 0L);
Message message = createJMSMessage(session, mail, props);
producer.send(message);
}
} catch (Exception e) {
throw new MailQueueException("Unable to enqueue mail " +
mail.getName(), e);
}
}
/**
* Dequeue ready messages.
*
* CRITICAL FIX FOR JAMES-4192:
* We DO NOT provide any message selector here!
*
* - The consumer receives messages directly without any selector filter.
* - Delayed messages are held by Artemis's ScheduledDeliveryHandler and
only
* released to the queue when their scheduled time arrives.
* - Head-of-line blocking and memory paging deadlocks are completely
eliminated.
*/
@Override
protected Mono<MailQueueItem> deQueueOneItem() {
Session session = null;
MessageConsumer consumer = null;
try {
session = connection.createSession(true,
Session.SESSION_TRANSACTED);
Queue queue = session.createQueue(queueName.asString());
// NO SELECTOR USED HERE:
// consumer = session.createConsumer(queue);
// Rather than:
// consumer = session.createConsumer(queue, getMessageSelector());
consumer = session.createConsumer(queue);
Message message = consumer.receive(10000);
if (message != null) {
return createMailQueueItem(session, consumer, message);
} else {
session.commit();
closeConsumer(consumer);
closeSession(session);
}
} catch (Exception e) {
rollback(session);
closeConsumer(consumer);
closeSession(session);
return Mono.error(new MailQueueException("Unable to dequeue next
message", e));
}
return Mono.empty();
}
/**
* Overridden to return null so that any base class callers do not inject
* the obsolete JAMES_NEXT_DELIVERY selector.
*/
@Override
protected String getMessageSelector() {
return null;
}
}
```
4. WHY THIS DEFINITIVELY FIXES JAMES-4192
--------------------------------------------------------------------------------
1. No Selector Evaluation in Paging:
Because the consumer selector is removed entirely, Artemis does not need
to evaluate SQL selector expressions against memory or disk buffers.
2. Scheduled Messages Stored in Scheduler Index, Not in Queue Pages:
In Artemis, messages with `_AMQ_SCHED_DELIVERY` or `deliveryDelay` are held
in an in-memory red-black tree (or scheduler journal) and are not considered
part of the active queue until triggered. Thus, they cannot saturate
paging memory or block subsequent non-delayed emails.
3. Paging Works as Designed:
When active ready messages exceed memory thresholds, Artemis pages them to
disk.
As consumers pull messages without filter rejections, new messages are read
smoothly from page files in sequential order.
4. Performance & Simplicity:
Compared to Apache Pulsar (which requires ZooKeeper, BookKeeper, separate
requeue topics and filter actors), the Artemis native scheduled delivery
leverages built-in JMS 2.0 semantics with zero extra infrastructure.
================================================================================
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]