SPQR 1.3.0: a production-ready system for horizontal scaling of PostgreSQL
github.com
github.com
Other than that, thanks to a lot of the recent work on connection handling and concurrency since PG 11, Postgres is getting better and better actually using these additional resources well: https://www.enterprisedb.com/blog/performance-comparison-maj...
It was not hard to calculate the current max traffic or estimate the traffic growth over the next 10 years.
I did a demo of the system running on my laptop (all of it) + Postgres handling 100x the current data without too much difficulty.
Still they went with the "scale" solution because it was the right design. (and of course the consultants and me got quite lot more work todo so made a good deal more money)
I think a lot of work goes into horizontal scaling which is necessary at a certain scale, but very few people actually get anywhere near that scale. It can be important to understand which things are needed at your scale and where you can simply buy some beefier hardware. I've been at places where people run a dozen sharded DB servers with each server having 16GB of RAM. Maybe that's resume-driven-development where someone wants to say they've done that.
I agree it adds alot of complexity to the problem, which is another cost.
I guess this would be another argument for pay-as-you go cloud-managed DBs, despite being more expensive than rolling your own.
* DB backups are now (much) faster.
* Smaller backups means faster restores which reduces your RTO (Recovery Time Objective)
* If you have a well architectured application a catastrophic DB failure will now only impact a portion of your userbase instead of all of them.
There are probably more good reasons but these are the ones I could think of now.(There are other reasons, of course...)
It usually scaled linearly. We had a 32-core (64-vcore) server, saturating all cores and running a bit more than 32x as fast as a single query. In some cases, it was less than linear but much better than singular, and I think that was only cause of mistakes like uuid4 pkeys.
Biggest limiter is memory, where the need for it grows linearly with table index size. Postgres really really wants to keep the index pages hot in the OS cache. Gets very sad and weird if it can’t: will unpredictably resort to table scans sometimes.
We are running on AWS Aurora, on a db.r6i.12xlarge. Nowhere even close to maxed out on potential vertical scaling.
Is it IoT / remote sensing related?
I worked for a company that had only a few thousand active customers yet had dozens of terabytes of data, for this reason.
Maybe not the best idea, I guess a file system would be better for that, and just use the DB for metadata.
But OTOH all the data is one place, so you just migrate the DB. Less to worry about.
I just looked up, all of English Wikipedia (including images) is barely even 100 GB ... crazy world we live in.
This one data store is easier is a myth , it just offloads complexity from developer to infra teams who are now provisioning premium NVMe storage instead of cold object stores for binary data .
Binary data is not indexed or aggregated in a SQL store there is no value in doing this is one place expect dev experience at the cost of infra team experience.
Timeseries (like IoT you mentioned ) or binary blobs or logs or any other data in SQL storage that shouldn’t be really there can hit any size wouldn’t be all that interesting.
Can’t speak for OP, however managing data for few million user apps, what I have observed is most SQL stores hit single TB range and then start getting broken down into smaller dbs either coz now teams have grown want their own Micro-service or DB or infra wants easier to handle in variety of ways including Backup /recovery larger DBs are extremely difficult to get reasonable RTO/RPO numbers for.
I think the reason behind aurora pick is to support arbitrary aggregation, filtering and low latency read (p90 < 3000ms). We could not pick distributed DB based on Presto, Athena or Redshift mainly for latency requirements.
The other contender I consider is Elastic search. But, I do think using it in this case is akin to fitting a square peg in round hole saying.
EDIT: Here's what I was thinking about. It's chunked in 10gb increments that are replicated across AZs.
> Fault-tolerant and self-healing storage
Aurora's database storage volume is segmented in 10 GiB chunks and replicated across three Availability Zones, with each Availability Zone persisting 2 copies of each write. Aurora storage is fault-tolerant, transparently handling the loss of up to two copies of data without affecting database write availability and up to three copies without affecting read availability. Aurora storage is also self-healing; data blocks and disks are continuously scanned for errors and replaced automatically.
If I understand correctly what you mean, then this is no longer a problem. You will simply need to use a connection pool, such as Odyssey or PgBouncer. Even SPQR has its own pool of connections for each shard.
I/O or compute constraints are another issue, if your CPU or disc is already saturated you get probably no additional benefit.
But if you wait for something (locks, I/O) the connection/process can't do other things. High latency between app/database and long running transactions can also use up your available processes, even if they don't consume a lot of CPU or I/O or fight for the same locks.
Lock contention is its own problem, but makes the blocking of processes/connections worse.
Partial indexes and partitioning go a loooooooong way.
In my experience if you want to be cost effective you need both. A decent amount of vertical scaling to have headroom for baseline and some amount of unpredictable spikes, horizontal scaling for the valleys of traffic that match your primary markets day/night cycle.
Repeat with bigger and bigger nides as needed. To scale down, do the inverse.
Like, Citus's FAQ says "if you use Citus, you do not need to manually shard your application, and you do not need to re-architect your application in order to scale out." But the line between application-level and DB-level sharding isn't this sharp. The fundamental limitations of distributed systems surface in their rules* about what you can/cannot do across shards, and you might find yourself re-architecting your application anyway.
* https://docs.citusdata.com/en/stable/develop/reference_worka...
https://github.com/pg-sharding/spqr/blob/1.3.0/LICENSE (BSD2)
https://github.com/citusdata/citus/blob/v12.1.2/LICENSE (AGPLv3)
> Notwithstanding any other provision of this License, if you modify the Program, your modified version must prominently offer all users interacting with it remotely through a computer network (if your version supports such interaction) an opportunity to receive the Corresponding Source of your version by providing access to the Corresponding Source from a network server at no charge, through some standard or customary means of facilitating copying of software. This Corresponding Source shall include the Corresponding Source for any work covered by version 3 of the GNU General Public License that is incorporated pursuant to the following paragraph.
In other words, if you deploy a modified version of Citus, your modified version is AGPL licensed, and thus you must provide the source code of this modified version to all users who interact with it remotely (e.g. through your web application).
What it does not state is that you must provide the source code of your entire web application just because you deployed a modified version of Citus, nor that your web application becomes an AGPL-licensed derived work of Citus because you used it over a network. The AGPL also does not require anything at all from you other than the plain GPL's basic requirements if the version of Citus you deploy is unmodified. (These are all incredibly common misconceptions on the internet, by people who've never read the license nor the GNU website.)
"The AGPLv3 does not adjust or expand the definition of conveying. Instead, it includes an additional right that if the program is expressly designed to accept user requests and send responses over a network, the user is entitled to receive the source code of the version being used."
Have I misread this? As I understand AGPL3 is an extension of GPL 3. Therefore means if I use an AGPL3 licensed service I have rights to receive the code for this service.
ref: https://www.fsf.org/bulletin/2021/fall/the-fundamentals-of-t...
That's why the text of the clause itself makes a distinction about "if you modify the Program".
https://www.asterix-obelix.nl/index.php?page=hjh/dos-italy.i...
https://www.youtube.com/watch?v=wjOfQfxmTLQ
It's funnier if you know that the ear-twisting Roman at the start was basically a caricature of many Latin teachers that were around when Latin was taught more widely in schools (such as when the Monty Python team were young).
Only azure supports Citus as far as I know.
(I expected Roman puns and was not disappointed)
I see nothing about partition tolerance, so I assume it isn't at all.
Give me Aphyr tests or there is no reason to pay attention.