This message was deleted.
# general
s
This message was deleted.
d
The short answer is no. MessageIds in Pulsar aren’t single counters, instead they are composed of multiple elements that indicate the batch id, entry id, ledger id, etc. This makes it impossible to consistently determine the message id of the current message id plus N because any of those elements may have changed if they are in a different batch or ledger, etc., e.g.
ledgerId:entryId:partitionIndex:batchIndex
j
The exact number of messages is hard to determine especially when an entry can consist of many messages and the entry data it self is compressed. Though I wonder if we can attach metadata to each entry that describes how many messages are in this entry. What we can determine is the size, so we can have an client API that asks for x bytes starting from a certain message id
πŸ‘ 2
@merlimat thoughts?
c
I notice that most pulsar source and sink tests are dropped in this commit https://github.com/streamnative/pulsar-spark/commit/c467f0996a0e5e338175ccd56db072facfe08e1e. Now the test coverage is insufficient, can we add it back?(or add other testing) @Neng @merlimat @jerry
n
yes, we should. But the previous test case may need update
c
Do you have plan to add the tests back in the near future?
n
given our current resource is tight (need to prepare various things for Pulsar Summit in Q3), we will probably start to work on that in Q4.
πŸ‘ 1
c
I can create a PR to add back the testing, could you elaborate on what update is necessary? @Neng
n
just make sure it follows some general scala test guideline.
we can refine them over the time.
really thank you for pick this up!
c
No problem! Can I simply add back all the old tests except the continuous one that is no longer supported?
n
sure
BTW, did you guys observe the spark connector creates too many connections on the pulsar topic issue from your customer?
c
Also, can someone point me to the pulsar client interface that we can leverage to implement admission control? I remember it was about polling the size of the next few batches of message and stop when the accumulated size is larger that read limit, but I don't know which interface I can use to achieve this. @Neng @merlimat @jerry
We have not heard from customer about any performance issue of pulsar spark connector. Is there any improvement you have in mind?
n
there might be a producer reuse optimization, but we are still investigating an issue reported by one of our users.
c
I saw that pulsar source RDD use preferred location to maximize reader reuse. I don't find similar optimization in the sink side.
n
yeah, this might be an issue
c
This is the implementation of cached kafka producer pool in spark, we can do something similar for pulsar sink https://github.com/apache/spark/blob/master/connector/kafka-0-10-sql/src/main/scal[…]che/spark/sql/kafka010/producer/InternalKafkaProducerPool.scala
kafka producer is thread safe, so it can be put in a cache and reused by multiple worker threads safely.
d
Can we prioritize this fix? We have a customer that is being impacted by this. @Neng @Chaoqin Li
πŸ‘€ 1
c
it seems that pulsar producer is reused in the executor, so the only problem is locality? https://github.com/streamnative/pulsar-spark/blob/master/src/main/scala/org/apache/spark/sql/pulsar/PulsarWriteTask.scala#L137 @Neng @David K
d
Thanks @Chaoqin Li, is there an eviction thread for this Map similar to the one for Kafka? We are seeing an ever increasing number of connections created to the Pulsar broker for the topic and are trying to determine the root cause
c
I don't see any eviction. Why would connection keeps increasing if the output topics are fixed?
d
Actually, if singleProducer is true, then a new producer get created every call, correct?
Copy code
if (null != singleProducer) {
      return singleProducer.asInstanceOf[Producer[T]]
    }
Copy code
protected lazy val singleProducer =
    if (topic.isDefined) {
      PulsarSinks.createProducer(clientConf, producerConf, topic.get, pulsarSchema)
    } else null
c
when is the topic not defined?
πŸ€” 1
d
I am not very scala fluent, so I could be wrong. But in the getProducer method the first call is to check if
singleProducer
is not null. If the topic IS defined, then a producer is created and returned.
This created/returned object is never added to the Map for reuse, or am I just missing something?
c
I think you are correct. A new PulsarWriteTask is created for each write and neither singleProducer or topic2Producer get reused.
We would need a singleton cache for this.
πŸ‘ 1
n
since write task may be scheduled to different executors, will a cache work in this case?
c
We don't enforce write tasks to be on the same executor explicitly, but they are usually scheduled on the same executor because we set locality for read task and stateful partition, and write tasks following them should also have certain level of locality.
πŸ‘€ 1
d
Should I create an issue to track this?
πŸ‘ 1
c
Can we have a simple fix before implementing the producer pool? Just call close() for the producer when the write task finished. @Neng @David K
πŸ‘€ 1
d
I guess that is better than the current scenario. But there would be a lot of network connections opened and then closed right away.
Why not just add the created producer to the map after it is created? Subsequent calls would then just reuse it rather than having to create a new one, use it, and then close a producer over and over again
c
The map is owned by the writer task, which is a short live object that doesn't get reused. We would need to implement a singleton cache with cache eviction, which would require some design and discussion.
j
πŸ‘
Probably can reuse much of the code for the consumer and producer cache pool from kafka source/sink
πŸ‘€ 1
πŸ‘ 1
n
https://github.com/streamnative/pulsar-spark/pull/137 <-- I draft a PR to fix the producer spam issue.
πŸ‘ 2
c
This fix looks good to me(at least no longer leak the connection), we can design and implement producer cache later.
πŸ‘ 1
We would also like to add pyspark support to pulsar spark connector, are you ok with that? @Neng @merlimat
πŸ‘ 1
n
sounds great to me!