Hi Jon,

NATS is basically a message queue, like ActiveMQ I suppose.

I included the adapter into the project using maven
<dependency>
     <groupId>io.nats</groupId>
     <artifactId>java-nats-streaming</artifactId>
     <version>2.2.3</version>
</dependency>

i started up a nats server using docker.  here is my docker-compose.yml
version: '3.1'
services:
  nats-docker:
    image: nats-streaming:0.17.0
    restart: always
    command:
      - '-p'
      - '4222'
      - '-m'
      - '8222'
      - '-hbi'
      - '5s'
      - '-hbt'
      - '5s'
      - '-hbf'
      - '2'
      - '-SD'
      - '-cid'
      - 'yourclientid'
    environment:
      TZ: Europe/London
      LANG: en_GB.UTF-8
      LANGUAGE: en_GB:en
      LC_ALL: en_GB.UTF-8
    ports:
      - '4222:4222'
      - '8222:8222'
    expose:
      - 4222
      - 8222
    networks:
      - backend
networks:
  backend:
    driver: bridge

JCA sounds good if it solves the threading issue.  it is very kind of you to offer to help write an adapter.  looking at the code you sent it looks complicated but i can have a stab at it if you don't have much time

let me know if you need more info

Matthew

On 07/06/2021 17:48, Jonathan Gallimore wrote:
At the risk of sounding a bit ignorant... what is NATS?

 From what I can tell, it sounds like you're receiving a stream of events
(over websocket) and want to do some processing in an EJB or CDI bean for
each event. The connection to the NATS server isn't in the context of a
HTTP (or any other type of) request, and just runs all the time while the
server is running - does that sound about right?

Assuming that sounds right, it sounds a bit like the Slack JCA connector I
wrote a while back:
https://github.com/apache/tomee-chatterbox/tree/master/chatterbox-slack.
Essentially, the resource adapter connects to slack and runs all the time.
Messages that come into the server from slack are processed in MDBs that
implement the InboundListener interface.

JCA certainly feels complex, especially when compared with your
Singleton @Startup bean approach, but I usually find that if I try and work
with threads in EJBs, things usually go in the wrong direction. Conversely,
JCA even gives you a work manager to potentially handle that stuff.

If you can give me some pointers to running a NATS server, I'd be happy to
help with a sample adapter and application.

Jon

On Mon, Jun 7, 2021 at 11:49 AM Matthew Broadhead
<[email protected]> wrote:

I am trying to subscribe to a NATS streaming server with
https://github.com/nats-io/stan.java which is java.lang.Autocloseable.

At first it wasn't closing properly as seen in my original gist:
https://gist.github.com/chongma/2a3ab451f2aeabc98340a9b897394cfe

This was solved with this

https://stackoverflow.com/questions/39080296/hazelcast-threads-prevent-tomee-from-stopping


creating a default producer:
@ApplicationScoped
public class NatsConnectionProducer {

      @Resource(name = "baseAddressNats")
      private String baseAddressNats;

      @Produces
      @ApplicationScoped
      public StreamingConnection instance() throws IOException,
InterruptedException {
          StreamingConnectionFactory cf = new
StreamingConnectionFactory(new Options.Builder().natsUrl(baseAddressNats)
.clusterId("cluster-id").clientId("client-id").build());
          return cf.createConnection();
      }

      public void destroy(@Disposes final StreamingConnection instance)
              throws IOException, TimeoutException, InterruptedException {
          instance.close();
      }
}

But now i am creating a new thread because any injections with JPA had
cacheing issues and this seems to work but i am not sure it is
broadcasting to websockets correctly
@Singleton
@Lock(LockType.READ)
@Startup
public class SchedulerEvents {
      private static final Logger log =
Logger.getLogger(SchedulerEvents.class.getName());

      @Inject
      private StreamingConnection streamingConnection;

      @Inject
      private SomeController someController;

      @PostConstruct
      private void construct() {
//        log.fine(Thread.currentThread().getName());
          try {
              streamingConnection.subscribe("scheduler:notify", new
MessageHandler() {
                  @Override
                  public void onMessage(Message m) {
                      try {
                          log.fine(Thread.currentThread().getName());
                          // this needs to spawn a new thread otherwise
injections are stale
                          Thread thread = new Thread(new Runnable() {
                              public void run() {
log.fine(Thread.currentThread().getName());
                                  process(m.getData());
                              }
                          });
                          thread.start();
                          while (thread.isAlive()) {
                              // wait
                          }
                          log.fine("Thread finished OK");
                          m.ack();
                      } catch (Exception e) {
                          emailController.emailStackTrace(e);
                      }
                  }
              }, new

SubscriptionOptions.Builder().startWithLastReceived().manualAcks().ackWait(Duration.ofSeconds(60))
                      .durableName("scheduler-service").build());
          } catch (IOException | InterruptedException | TimeoutException e)
{
              e.printStackTrace();
          }
      }

      private void process(byte[] data) {
          String raw = new String(data);
          JsonReader jsonReader = Json.createReader(new StringReader(raw));
          JsonObject jo = jsonReader.readObject();
          jsonReader.close();
          String type = utilityDao.readJsonString(jo, "type");
          int id = utilityDao.readJsonInteger(jo, "id");
          if (type == null || id == 0) {
              emailController.emailThrowable(new Throwable(), raw);
              return;
          }
          log.info("Received a message: id: " + id + ", type:" + type);
         DefaultServerEndpointConfigurator dsec = new
DefaultServerEndpointConfigurator();
         SomeWebSocket nws = dsec.getEndpointInstance(SomeWebSocket.class);
         nws.broadcast(ja.toString());
      }

}

what is the best way to use an autocloseable?




Reply via email to