Hi ! we use the pulsar connector to feed several r...
# pinot-dev
m
Hi ! we use the pulsar connector to feed several realtime tables on Pinot. We find that the connector creates a lot of subscriptions on Pulsar which are never cleaned (hundreds or thousands per day per topic, depending on the events flow). We are forced to configure Pulsar to clean up unused subscriptions quickly, but I don't think that's a good solution. Is this behavior normal? Thank you !
n
@Mathieu Druart just to clarify: are you using pulsar with low level consumer or high level consumer (as in your table config) ?
m
@Navina low level
n
ok. so I am fairly new to pulsar. Is my understanding right - that a new subscription gets created every time a pulsar client is instantiated?
m
A subscription is created by a Consumer (not by the Client). A subscription is persistent and designed to be reused over time to track messages consumption (even in case of disconnection/reconnection). When creating a Consumer if the subscription identifier is reused, message consumption will reuse/continue the same subscription.
👍 1
as a subscription is persistent, in theory it should be reused or else deleted if it is no longer useful.
n
@Mathieu Druart I am aware of some code paths in pinot that re-create consumers periodically due to idleTimeout in the stream (can happen with low volume topics). Additionally, due to some legacy reasons, the metadata provider instantiates a new consumer every time it tries to fetch offsets or other metadata. See here and here. These need to fixed, hopefully soon 🤞
In the short term though, can you try increasing the idleTimeout for your tables to a very high value and see if it reduces the number of subscriptions (with random id) it generates?
m
thank you @Navina We will try this and let you know
n
cool . yeah these refactoring have been bugging me for a while now. Will try to improve it soon.
🙏 1
@Mathieu Druart qq: do we need to make an explicit call to delete a subscription ? In that case, I don't think we delete it anywhere in our codebase 🙈
m
@Navina Maybe I don't know all the logic on the Pinot side, but normally the subscription must be reused over time (it's the subscription that follows the progress with the offset of the last message consumed). If it is not possible to reuse the same subscription, it probably depends on the quantity of new subscription created but it would be better to delete them yes. On our side, we have configured Pulsar to delete unused subscriptions after 2 days.
what is a bit strange is that if each time a Consumer is created, a new subscription is created, it looks more like the use case of a Reader rather than a Consumer
n
I see. Yeah Pinot manages its own offsets for consumption and hence, uses the
Reader
api in pulsar. It will seek to a specific messageId each time it fetches. So, I don't think the subscriptions are created when the pinot partition consumer is created.
the consumer that is being created is in the pinot's stream metadata provider. It creates a "consumer" rather than a "reader" to fetch partition metadata . I think we need to add clean up in that part of the code.
@Mathieu Druart tks for clarifying. Looks like in Pulsar's case, we can trivially solve it by deleting the subscription explicitly after metadata provider is closed.
what is a bit strange is that if each time a Consumer is created, a new subscription is created, it looks more like the use case of a Reader rather than a Consumer
you are right. we are using the reader api.
m
thank you @Navina
👍 1
just a question @Navina the consumer is only used to get the last message id ?
you can get the last message id with the Pulsar admin I think, without creating Consumer
Copy code
PulsarAdmin.builder().serviceHttpUrl("HTTP URL to Pulsar Admin API").build().topics().getLastMessageId(topic);
but I don't know if it covers all the needs or if the performances are equivalent
n
@Mathieu Druart we also want to know the first message ID if the reset policy is earliest. Tks for the reference. I have mentioned this in the PR - https://github.com/apache/pinot/issues/9854
🙏 1
I think there is an
examineMessage
API that can help fetch both offsets. I need to try it out to see if it can be used here instead of spinning up a consumer
👍 1
m
just to let you know the examineMessage API is only available in Pulsar 2.8.1 and later versions. The Pulsar client of Pinot is in version 2.7.2.
n
oh good to know 🙂
m
Hi @Navina increasing the idle timeout doesn't seem to change anything, which doesn't surprise me because it's a topic with a fairly high flow all the time (I don't think several seconds go by without a message).
n
Ok got it. I have a PR to delete the subscriptions that are getting created. Once it gets merged, please try it out. It might alleviate your problem
m
Hi @Navina I added a comment on your PR
👀 1