This message was deleted.
# general
s
This message was deleted.
m
I was using v1.7.0. I also noticed a new version of the pulsar client, v1.8.0, came out 4 days ago. Which is great. So I also tried this out. Same results. Around 10 messages are being retrieved, after that it just stops.
Copy code
const Pulsar = require('pulsar-client');

(async () => {
  const client = new Pulsar.Client({
    serviceUrl: '<pulsar://pulsar:6650>',
  });

  const reader = await client.createReader({
    topic: '<persistent://public/default/test-topic>',
    startMessageId: Pulsar.MessageId.latest(),
    listener: (msg,reader)=>{
      const msgData = msg.getData().toString();
      console.log(msgData)
    }
  });
})();
Bytheway: I'm having loads of fun with Pulsar..... My new technology for 2023 🎉
After debugging a while, it seems when the last record is received, the pulsar client closes it's connection. I think this is by design, but hope someone can confirm. So I'm sharing a bit of my learning experience here. Hope it helps someone. This is the code I use now. So for every topic I have a separate process:
Copy code
const Pulsar = require('pulsar-client');

let reader;

(async () => {
  const client = new Pulsar.Client({
    serviceUrl: '<pulsar://pulsar:6650>'
  });

  reader = await client.createReader({
    topic: '<persistent://public/default/test-topic>',
    startMessageId: Pulsar.MessageId.latest(),
    // listener: (msg,reader)=>{
    //   const data = msg.getData();
    //   console.log(msgData.toString())
    // }

  });

  function readNext() {
    return new Promise((resolve, reject) => {
      setTimeout(async () => {

        const msg = await reader.readNext();
        const msgData = msg.getData();

        console.log(msgData.toString());
        
        resolve()
      }, 10)
    })
  }

  while (true){
    await readNext();
  }
})();
This feels a bit like an anti-pattern...
d
I would recommend using the consumer API if your intent is to continuously read messages from the topic as the arrive. I think your issue might have been due to the fact that you weren’t acknowledging the messages after you consumed them, which can lead to back-pressure
Copy code
(async () => {
  // Create a client
  const client = new Pulsar.Client({
    serviceUrl: '<pulsar://localhost:6650>',
  });

  // Create a consumer
  const consumer = await client.subscribe({
    topic: 'my-topic',
    subscription: 'my-subscription',
    subscriptionType: 'Exclusive',
  });

  // Receive messages
  while (true) {
    const msg = await consumer.receive();
    console.log(msg.getData().toString());
    consumer.acknowledge(msg);
  }

  await consumer.close();
  await client.close();
})();
or with a message listener
Copy code
// Create a consumer
const consumer = await client.subscribe({
  topic: 'my-topic',
  subscription: 'my-subscription',
  subscriptionType: 'Exclusive',
  listener: (msg, msgConsumer) => {
    console.log(msg.getData().toString());
    msgConsumer.acknowledge(msg);
  },
});
m
I didn't try the ack yet, if this works it means I can get rid of the endless loop.... will give it a try now!
Tried it. I can confirm the acknowledge in the listener does not really prevent the client from closing the connection after it catches up with the last message in the topic.
🤔 1
I'm now indeed leaving the reader approach, subscribing to multiple topics. (This only is possible with consumers). And then I use the msg.getTopicName() to differentiate and send the message to a nodejes EventEmitter, so it can be picked up by the frontend. (I use server sent events now, those really are nice to work with).
Could be it is the NodeJS version. Wild guess..
d
yea, that seems odd that the client would close the connection. Does your code tear down the client object itself? Have to ask …
m
Nope.
So I put the client outside the main method, I attached the debugger, and after the last message came in, I discovered the client was in a closed state.