Prometheus uses the pull principle. What does that mean, and what is its advantage when you run several Prometheus instances?
The server keeps a list of targets and fetches ("scrapes") their metrics in turn; targets don't send anything on their own. So targets need no configuration about where to send data. Only the Prometheus instances decide whom they ask.
* Adding a second server: with pull only the new server changes; with push every host does. *
With push-based monitoring, every monitored host has to be configured with the address of the server it reports to. Add a second server for high availability, or split the load across instances, and you have to reconfigure every host.
With pull, the targets just offer their metrics on an HTTP endpoint and don't care who reads them. Running two Prometheus instances for redundancy is simply two servers scraping the same targets; sharding is giving each instance a different subset of targets. Nothing changes on the hosts.
A bonus: if a scrape fails, Prometheus knows immediately that the target is unreachable. The synthetic up metric becomes 0, which is the basis of the classic "instance down" alert.
Go deeper:
Prometheus blog — Pull doesn't scale, or does it? — the developers' own case for pull at scale.
Prometheus FAQ — includes "Why do you pull rather than push?".