This message was deleted.
# general
s
This message was deleted.
m
The problem does not occur when there are no writes/reads.
I’ve created a simple chaos test that can be run on clean Pulsar: https://github.com/websight-io/pulsar-chaos-test Current behaviour: • messages are lost • some partitions are deleted
a
@Penghui You’ve done those, right?
m
Please note, that the tests are basic, just to reproduce the error. It’s critical for us, we cannot accept any data loss. Should I create a GH issue for that ?
p
@Michał Cukierman Do you have subscriptions or a data retention policy for the topic? If no subscriptions and no data retention, Pulsar will remove the data automatically.
One thing you can try is to create a subscription on the topic, then start the produce program, then restart all the brokers.
s
@Michał Cukierman: Add more to what Penghui said: 1. You should try to run consumer before run producer. So the consumer will create a subscription to hold the messages. 2. You also need to configure the retention policy to ensure the data is kept even they are ackwnoledged. 3. https://github.com/websight-io/pulsar-chaos-test/blob/main/consumer_verify.py#L4 <- for consumer to verify the data, you want to specify the subscription initial position to be earliest, so your verify consumer can reply the messages from the beginning.
m
The topic is created from before any messages are produced with retention set to
-1
: https://github.com/websight-io/pulsar-chaos-test/blob/main/run_prepare_topic.sh#L88
bin/pulsar-admin topics set-retention persistent://${tenant}/${namespace}/monolog --size -1 --time -1
Consumer is created after the producer, but it’s possible to read all the messages including `verification-message`: https://github.com/websight-io/pulsar-chaos-test/blob/main/run.sh#L21C1-L21C42 I am receiving:
Copy code
"OK 'verification-message' was present in 'monolog' topic before restart"
from the: https://github.com/websight-io/pulsar-chaos-test/blob/main/run_verify.sh#L18 So yes, all the messages and the subscription were on the topic while reading. And retention worked before the restart. @sijieg I’ll correct that. I don’t think that the result will change, since all of the partitions are removed after restart I get:
Partitions need to create
X
In general I can also see that the partitions are removed from `pulsar-manager`` , and I see the ledgers are removed ((
Grafana
))
I’ll correct the script to run consumer before the producer (this should not affect anything, since the retention is set) I’ll also set the initial position to
Earliest
, but again I should not loose the ledgers and data
In a meantime I am stucked with another error. Cannot remove the partition using command:
bin/pulsar-admin topics delete-partitioned-topic <persistent://websight/dxp/monolog> --force
I also removed the subscriptions, but it didn’t help. I get no errors in the log files.
I will run the tests later today. It looks like I lost my window and the environment I’ve provisioned is already shutting down (company policy)
Changes are pushed to the repository
I’ve run the test again. 1. Messages received before restart:
Copy code
Received message: 'b'verification-message''
Received message: 'b'hello-pulsar-0''
Received message: 'b'hello-pulsar-1''
Received message: 'b'hello-pulsar-2''
Received message: 'b'hello-pulsar-3''
Received message: 'b'hello-pulsar-4''
Received message: 'b'hello-pulsar-5''
Received message: 'b'hello-pulsar-6''
Received message: 'b'hello-pulsar-7''
Received message: 'b'hello-pulsar-8''
Received message: 'b'hello-pulsar-9''
Received message: 'b'hello-pulsar-10''
...
2. Messages received after restart:
Copy code
Received message: 'b'hello-pulsar-0''
Received message: 'b'hello-pulsar-4''
Received message: 'b'hello-pulsar-12''
Received message: 'b'hello-pulsar-24''
Received message: 'b'hello-pulsar-36''
Received message: 'b'hello-pulsar-48''
Received message: 'b'hello-pulsar-60''
Received message: 'b'hello-pulsar-72''
Received message: 'b'hello-pulsar-84''
Received message: 'b'hello-pulsar-96''
Received message: 'b'hello-pulsar-108''
Received message: 'b'hello-pulsar-120''
Received message: 'b'hello-pulsar-132''
Received message: 'b'hello-pulsar-144''
Received message: 'b'hello-pulsar-156''
Received message: 'b'hello-pulsar-168''
Received message: 'b'hello-pulsar-180''
Received message: 'b'hello-pulsar-192''
Received message: 'b'hello-pulsar-204''
Received message: 'b'hello-pulsar-216''
Received message: 'b'hello-pulsar-228''
Received message: 'b'hello-pulsar-240''
Received message: 'b'hello-pulsar-252''
Received message: 'b'hello-pulsar-264''
Received message: 'b'hello-pulsar-276''
Received message: 'b'hello-pulsar-288''
Received message: 'b'hello-pulsar-300''
Received message: 'b'hello-pulsar-312''
Received message: 'b'hello-pulsar-324''
Received message: 'b'hello-pulsar-336''
Received message: 'b'hello-pulsar-348''
Received message: 'b'hello-pulsar-360''
Received message: 'b'hello-pulsar-372''
Received message: 'b'hello-pulsar-384''
Received message: 'b'hello-pulsar-16''
Received message: 'b'hello-pulsar-28''
Received message: 'b'hello-pulsar-40''
Received message: 'b'hello-pulsar-52''
Received message: 'b'hello-pulsar-64''
Received message: 'b'hello-pulsar-76''
Received message: 'b'hello-pulsar-88''
Received message: 'b'hello-pulsar-100''
Received message: 'b'hello-pulsar-112''
Received message: 'b'hello-pulsar-124''
Received message: 'b'hello-pulsar-136''
Received message: 'b'hello-pulsar-148''
Received message: 'b'hello-pulsar-160''
Received message: 'b'hello-pulsar-172''
Received message: 'b'hello-pulsar-184''
Received message: 'b'hello-pulsar-196''
Received message: 'b'hello-pulsar-208''
Received message: 'b'hello-pulsar-220''
Received message: 'b'hello-pulsar-232''
Received message: 'b'hello-pulsar-244''
Received message: 'b'hello-pulsar-256''
Received message: 'b'hello-pulsar-268''
Received message: 'b'hello-pulsar-280''
Received message: 'b'hello-pulsar-292''
Received message: 'b'hello-pulsar-304''
Received message: 'b'hello-pulsar-316''
Received message: 'b'hello-pulsar-328''
Received message: 'b'hello-pulsar-340''
Received message: 'b'hello-pulsar-352''
Received message: 'b'hello-pulsar-364''
Received message: 'b'hello-pulsar-376''
Received message: 'b'hello-pulsar-388''
Received message: 'b'hello-pulsar-8''
Received message: 'b'hello-pulsar-20''
Received message: 'b'hello-pulsar-32''
Received message: 'b'hello-pulsar-44''
Received message: 'b'hello-pulsar-56''
Received message: 'b'hello-pulsar-68''
Received message: 'b'hello-pulsar-80''
Received message: 'b'hello-pulsar-92''
Received message: 'b'hello-pulsar-104''
Received message: 'b'hello-pulsar-116''
Received message: 'b'hello-pulsar-128''
Received message: 'b'hello-pulsar-140''
Received message: 'b'hello-pulsar-152''
Received message: 'b'hello-pulsar-164''
Received message: 'b'hello-pulsar-176''
Received message: 'b'hello-pulsar-188''
Received message: 'b'hello-pulsar-200''
Received message: 'b'hello-pulsar-212''
Received message: 'b'hello-pulsar-224''
Received message: 'b'hello-pulsar-236''
Received message: 'b'hello-pulsar-248''
Received message: 'b'hello-pulsar-260''
Received message: 'b'hello-pulsar-272''
Received message: 'b'hello-pulsar-284''
Received message: 'b'hello-pulsar-296''
Received message: 'b'hello-pulsar-308''
Received message: 'b'hello-pulsar-320''
Received message: 'b'hello-pulsar-332''
Received message: 'b'hello-pulsar-344''
Received message: 'b'hello-pulsar-356''
Received message: 'b'hello-pulsar-368''
Received message: 'b'hello-pulsar-380''
Received message: 'b'hello-pulsar-392''
Traceback (most recent call last):
  File "/Users/michalcukierman/dev/pulsar/chaos_test/consumer_verify.py", line 9, in <module>
    msg = consumer.receive(5000)
          ^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.11/site-packages/pulsar/__init__.py", line 1277, in receive
    msg = self._consumer.receive(timeout_millis)
          ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
_pulsar.Timeout: Pulsar error: TimeOut
Some other errors: There is backlog on subscription, but all the messages are received
Some partitions are empty after restart:
Output of
bin/pulsar-admin topics get-retention <persistent://websight/dxp/chaos-test>
Copy code
{
  "retentionTimeInMinutes": -1,
  "retentionSizeInMB": -1
}
The situation looks OK when I am restarting up to 2 brokers:
Copy code
➜  out git:(main) ✗ cat consumer.log | grep "hello-pulsar" | wc -l    
    1420
➜  out git:(main) ✗ cat verification.log | grep "hello-pulsar" | wc -l
    1460
I’ve got more messages during the verification, but that may happen. The data loss is happening when I am restarting 3 brokers out of 3 at the same time.
Partitions look OK
IT looks like all the entries are in bookies, but are not assigned to the topic.
From interesting (I guess) warnings when the data loss occur, I see:
Data loss during brokers restart is happening on local setup as well. Is it how Pulsar is suppose to work? The tests above are run on the GCP. Today, I’ve created local cluster running on MicroK8S, the configuration of helm values is here: https://github.com/websight-io/pulsar-chaos-test/tree/main/environment-local (deploy.sh, and values.yaml) The configuration is: 1 BK, 3 brokers, 1 proxy, 1 ZK After the restart of 3 brokers during reading/writing I’ve lost more than 50% of the messages from the topic:
Copy code
Results:
NOK Missing partitions found
OK 'verification-message' was present in 'chaos-test' topic before restart
OK 'verification-message' is present in 'chaos-test' topic after restart
Messages available before restart:
    2614
Messages available after restart:
    1160
The retention is set:
Copy code
{
  "retentionTimeInMinutes": -1,
  "retentionSizeInMB": -1
}
I am reading the topics using
Earliest
position. The whole scenario can be easily reproduced using given repository. There is no external dependencies other than Pulsar 3.0.
s
@Michał Cukierman You are using non-persistence storage.
Copy code
volumes:
  persistence: false
In Kubernetes, it will uses emptyDir. So once the bookie or zk pod is restarted, the data is gone.
m
on GCP I tested it with persistence: true, and the messages are still lost
the data are on bookies, which are not restarted
@sijieg the data loss is caused by the brokers restart
Plus there are 1160 messages available after restart, so the persistence works
I used it to simplify the local setup, and not to use PVC. Still it’s enough to reproduce the error
s
Broker is stateless. Restarting brokers will not result in data loss. I feel that there is a problem in the way how you have verified the results. Can you share the log file of
verification.log
,
consumer.log
?
m
@sijieg it’s commited to the repository
I know that the broker is stateless, that’s why I think it’s a bug
the backlog size for the new consumer decreases, plus you can observe that some partitions are empty after restart.
I’ll report a bug on GH, probably it will be easier to take it from there
Note that the repository I’ve created is just a projact to ilustrate the bug. We are loosing the messages in our product setup (much different from what’s on GH)
s
I believe your verification approach is problematic. If you generate the list of messages you have consumed, you will find the verify consumer received all the messages. Here are two files that I generated from your logs:
Copy code
cat consumer.log | grep "Received message" | grep "hello-pulsar" | awk '{ print $3 }' | sort | uniq > consume_messages.txt
Copy code
cat verification.log | grep "Received message" | grep "hello-pulsar" | awk '{ print $3 }' | sort | uniq > unique_messages.txt
if you compare the results, you will see will have both consumers received the messages up to
hello-pulsar-6326
and the verify consumer received the messages up to
hello-pulsar-6334
(I assumed the consumer was killed later).
From your logs, there is no data loss.
m
I might done a mistake here. But I for sure was loosing the messages - see the screenshots from Pulsar manager for example
So this one may be a successful run
I’ve just updated the Pulsar to 3.0.1 I hoped that it’s fixed here.
s
well. I don't think so. the scripts you wrote doesn't do query position by timestamp. Also, I will also question about Pulsar Manager's metrics. Because Pulsar Manager caches metrics in its own database, it doesn't use the metrics/stats data from Pulsar.
If you want to investigate any data loss issues, getting stats/stats-internal is the best to verify
m
Anyway, I still get: cat out/consumer.log | grep “hello-pulsar” | uniq | wc -l cat out/verification.log | grep “hello-pulsar” | uniq | wc -l I get: before: 9008 after: 6335
s
not pulsar manager
m
OK I get it, a lot of not acked re-deliveries
s
cat consumer.log | grep "Received message" | grep "hello-pulsar" | sort | uniq | wc -l 6327
you need to sort it then uniq
m
on 3.0.0. I was loosing the
verirication-message
messages as well
I’ve created it to make sure the one is not lost
anyway I need to do more tests on 3.0.1
It might have been fixed
s
I don't think there was an issue
m
Let me then get back to 3.0.0. I’ll try to update the repo with the results from 3.0.0
Thanks for your help!
s
You are welcome. If you have a problem in your production, feel free to raise a ticket to describe the symptom, we can take a look.
🙏 1
m
@sijieg I think I’ve reproduced it. I have the following results:
Copy code
Results:
NOK Missing partitions found
OK 'verification-message' was present in 'chaos-topic' topic before restart
NOK 'verification-message' is not present in 'chaos-topic' topic after restart
Messages available before restart:
     205
Messages available after restart:
      88
I am using your scripts to sort unique keys. I am now able to connect to the partition using new subscription and I always receive 88 items.
I’ve commited the logs to the repository
I use the new Python client:
Copy code
import pulsar

client = pulsar.Client('<pulsar://localhost:6650>')
consumer = client.subscribe('websight/dxp/chaos-topic',
                            subscription_name='NEW-SUBSCRIPTION_HERE',
                            initial_position=pulsar.InitialPosition.Earliest)

while True:
    msg = consumer.receive(5000)
    print("Received message: '%s'" % msg.data())
    consumer.acknowledge(msg)

client.close()
It was available before the restart, because the first consumer read it
What I did is to decrease the waiting time from 300 s to 30 s
Here is the video from the test run (I had to restart the proxy a couple of times, because the pods were not ready) Looks like early reads (starting verification early) may cause it?
And now I have the warnings:
[NEW-SUBSCRIPTION_HERE-777] Current position 209:-1 is ahead of last position 155:29"
What may suggest that the cursor is recovered in wrong state. I’ve checked, the messages are still in Bookie (checked in Grafana before)
s
It seems the verify consumer was killed before it is able to finish reading. How did you kill the verify consumer?
Copy code
Traceback (most recent call last):
  File "/Users/michalcukierman/IdeaProjects/websight-io/chaos_test/consumer_verify.py", line 9, in <module>
    msg = consumer.receive(5000)
          ^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.11/site-packages/pulsar/__init__.py", line 1277, in receive
    msg = self._consumer.receive(timeout_millis)
          ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
_pulsar.Timeout: Pulsar error: TimeOut
Can you try to change
receive(5000)
to
receive()
?
Copy code
while True:
    msg = consumer.receive(5000)
    print("Received message: '%s'" % msg.data())
    consumer.acknowledge(msg)
m
Sure I’ll record it, probably the program will wait and receive no messages
How long should I wait before terminating ?
I can run a producer and send one more message
s
That works
m
one arrived
It means that some previous messages are not accessible by the cursor, that would be maybe:
[NEW-SUBSCRIPTION_HERE-777] Current position 209:-1 is ahead of last position 155:29"
I’ve sent two, one arrived
but I’ve closed the producer just afterwards
s
Can you get topic stats/stats-internal for that topic and partitions?
m
doing it now
you mind giving me the command ?
I am using pulsar admin
got it
bin/pulsar-admin topics partitioned-stats-internal websight/dxp/chaos-topic
y
have you tried kill -15 instead of -9 for your producer/consumer? -9 won’t give process a chance to flush data in memory before exit.
also, it will be great if you share the purpose of chaos test. pulsar is a distributed system and shuting down all doesn’t make much sense in real cloud env.
m
This is happening in our product. I’ve wrote the chaos test to reproduce the issue on clean Pulsar. Shutting down all brokers happen, i.e. if you run out of memory after producing to much load. A bug in JVM, during wrong deployment or miscondigurarion. It may happen during migration. A bug in K8S release. Human error. I see plenty of things that may go wrong
The other thing is that you don’t know what causes the bug, and it may happen in other scenarious as well.
Regarding -9 vs -15 - The way on how someone closes the client should not affect the platform durability.
y
I agreed regarding those potential issues may happen. but those issues should be handled in different layers. those go back to what’s the possibility and how much effort you want pulsar to handle all of those scenarios.
pulsar-client batch messages in memory for retry in case broker is not available. if you use -15, your pulsar client will try to resend messsages in buffer. This is a part of pulsar handling failures.
The point here is that what you are testing makes sense (normal) and what next you want to do to mitigate. for example, jvm bug, you can swap/roll over a subset of broker and namespace isolation to segregate the issue. you can also use service mesh to do the same. it can be done in many different way, which one makes sense to you? it really depends.
more example, when you delete all broker pods at once, what scenario do you mimic? a k8s node down? a rack down? a zone/region down? or a whole public cloud down? In different scenarios, the speed to respawn new broker pods will be different. this goes back to what the chance of the whole public cloud down? how to mitigate that? how critical is your data to mitigate that?
the default pulsar has a script of rack awareness based on dns/subnet to get bookie rack awareness, but no default cloud zone/region awareness. you may need operator to handle the zone/region to avoid scenario like zone down.
m
I am aware. There is an issue, we observed it. We may report it and work on a fix or live with the bug. I prefer the first option. You don’t understand the cause of the bug until you spot it and fix it. Is that ok too? Loss of the data during restarts of ‘stateless’ services is critical our case, no matter of how likely it is. Its a risk we need to mitigate. Closing of the client does not change anything. You loose the historical data. Killing is done at the end. Still the data should be durable.
But again, we addressed it in our setup, I created a ticked in GH for it. People will be aware of the risk until its fixed. I am happy to support anyone who can fix it or to work on the PR sometime in the future. The latter would require some prior research as I am new to Pulsar. Still we are in the early stage of the product development lifecycle, we can pretty easy switch to other messaging system, if we encounter more issues like this.
👍 1