chovy-3012 opened a new issue, #12694:
URL: https://github.com/apache/seatunnel/issues/12694

   ### Search before asking
   - [x] I had searched in the 
[feature](https://github.com/apache/seatunnel/issues?q=is%3Aissue+label%3A%22Feature%22)
 and found no similar feature requirement.
   
   ### Description
   
   The ActiveMQ sink connector (`connector-activemq`) has a resource management 
bug in `write()` and lacks commonly used configuration options for the 
`ActiveMQConnectionFactory` and JMS `MessageProducer`.
   
   **1. Resource leak — Session and MessageProducer created per `write()` call**
   
   In the original `write()` method, a new `Session` and `MessageProducer` are 
created for every message and never closed. Under sustained throughput this 
leaks JMS resources (threads, sockets) and degrades broker performance.
   
   ```java
   // Original write() — creates new Session + Producer per call, never closed
   public void write(byte[] msg) {
       this.connection.start();
       Session session = this.connection.createSession(false, 
Session.AUTO_ACKNOWLEDGE);
       Destination destination = session.createQueue(config.get(QUEUE_NAME));
       MessageProducer producer = session.createProducer(destination);
       ...
   }
   ```
   
   **2. Missing commonly used config options**
   
   The connector exposes 8 optional `ActiveMQConnectionFactory` properties but 
omits several high-frequency options that users typically need to tune:
   
   - **ConnectionFactory-level**: `maxThreadPoolSize`, `sendTimeout`, 
`useCompression`, `connectResponseTimeout`, `producerWindowSize`, `useAsyncSend`
   - **Producer-level**: `deliveryMode`, `timeToLive`, `priority`
   
   Proposed solution:
   - Promote `Session` and `MessageProducer` to final fields created once in 
the constructor and reused in `write()`
   - Add 9 config options with defaults matching the ActiveMQ client library:
   
   | Option | Type | Default | Level |
   |--------|------|---------|-------|
   | `max_thread_pool_size` | int | 1000 | ConnectionFactory |
   | `send_timeout` | int | 0 | ConnectionFactory |
   | `use_compression` | boolean | false | ConnectionFactory |
   | `connect_response_timeout` | int | 0 | ConnectionFactory |
   | `producer_window_size` | int | 0 | ConnectionFactory |
   | `use_async_send` | boolean | false | ConnectionFactory |
   | `delivery_mode` | int | 2 (PERSISTENT) | Producer |
   | `time_to_live` | int | 0 (never expires) | Producer |
   | `priority` | int | 4 | Producer |
   
   ### Usage Scenario
   
   When using the ActiveMQ sink for high-throughput message delivery, users 
need to:
   - Avoid resource leaks from per-message Session/Producer creation
   - Tune delivery mode (PERSISTENT vs NON_PERSISTENT), message TTL, and 
priority
   - Control async send behavior and flow control for throughput vs reliability 
trade-offs
   - Configure connection timeout and thread pool sizing for the broker 
connection
   
   ### Related issues
   
   None
   
   ### Are you willing to submit a PR?
   
   - [x] Yes I am willing to submit a PR!
   
   ### Code of Conduct
   
   - [x] I agree to follow this project's [Code of 
Conduct](https://www.apache.org/foundation/policies/conduct)
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to