Postgres Count Performance
citusdata.com
citusdata.com
We use something similar to the trigger-based method they describe, tho have found that a lot of updates to count table inevitably ends with deadlocks. So instead of updating a count value, we always insert a new count of 1 or -1, and use summing to calculate the total count as needed. A background task is responsible for continually squashing the count values.
At this point, I wonder: with a huge dataset, wouldn't you have better time leaving the count field out of postgres altogether and use something like redis with its INCR/DECR instructions instead? This would prevent having deadlocks as well.
EDIT: that is, if you don't need to use the count field in other queries.
Another issue with the global tally is that the row_counts table bloats quite a bit, because the trigger is executed per-row and the update creates a copy of the row (so the next invocation has to walk all the previous MVCC copies, causing the long INSERT time). I wonder whether statement-level triggers might be used here, somehow (I don't think so).
> Note that custom extensions written in C like count_distinct are not bound by the value of work_mem. The array constructed in this extension can exceed your memory expectations.
This is slightly inaccurate, as this is not specific to custom aggregates. Custom aggregates are estimated just like any other (built-in) aggregates - number of expected groups times memory per group. The trouble is that (a) for aggregates with variable-length per-group state, we don't have a good size estimate, and (b) if the planner decides to use HashAggregate, we're unable to do anything when reaching work_mem. But this has nothing to do with the aggregate being custom - array_agg() and string_agg() have the same issue, for example.
FWIW, the extension was written quite a long time ago - before the various sort optimizations made by Peter Geoghegan. I wonder whether that made count_distinct obsolete.
Can you say a bit more why?
Jokes aside, it seems simple but is fairly tricky, as it depends both on input data and various other parameters. For example for array_agg() or string_agg() it might be estimated from number of entries / average length. For hll it also depends on the accuracy and expected number of values to track, etc. I was considering adding another method to the API, providing a better estimate, but never got to that.
So the current code simply assumes 1KB (IIRC) per group in those cases, or something like that.
Of course, the memory estimate also depends on the number of groups, but that's an orthogonal issue.
https://www.postgresql.org/docs/current/static/parallel-plan...
While doing the benchmarks, we could see that citus was always taking full advantage of all the cores in the cluster, while postgres parallelization was not.
Disclaimer: I'm not a db expert and I don't have any relationship with either citus or postgres.
I expect that's the same for all function calls like this? Surely pg has a concept of constants and doesn't needlessly re-check parameters?
Past the recommendation of counting based on an indexed column, I wonder if this should really be user's concern. This paragraph especially triggers a "this should be fixed upstream" feeling in me:
> A word of warning. When work_mem is high enough to hold the whole relation PostgreSQL will choose HashAggregate even when an index exists. Paradoxically, giving the database more memory resources can lead to a worse plan. You can force the index-only scan by setting SET enable_hashagg=false; but remember to set it true again afterward or other query plans will get messed up.
But worst case scenario, this article will be useful until this is fixed, so thanks again :)
Query planning is something where a poor choice can have serious performance ramifications, because n is usually much larger than in most programs. Analyzing the algorithmic complexity of a piece of SQL takes some experience and experimentation, and with different table stats the query planner may make different decisions. It can be worthwhile limiting the planner's discretion to get more predictable performance.
(I work on a product where many of the features can be expressed in terms of relational algebra. Often, both the best performing and quickest to write and test implementation logic is a bunch of SQL, and not the kind of CRUD that is easily wrapped with an ORM. What would make my life easier is a SQL linter that, given a model of costs, would prevent people writing queries that scale only linearly with specific table sizes. I've accumulated sufficient intuition that I could do this for MySQL at this point.)
But actually, I've just made a test, and it appears changing this setting only impacts the current connection, so provided it's toggled back after the request, this should not be a problem.
It would be interesting to see how much the performances improve once you use cstore_fdw (especially since 1M records is quite small when talking about OLAP workloads).
disclaimer: I've never used cstore_fdw, but I have evaluated a number of columnar databases in the past.
We find that the primary motivation for using cstore is reducing disk I/O / storage footprint. cstore_fdw keeps a columnar layout on disk in compressed form and reads only relevant columns. For example, it's commonly used for data archival purposes.
That said, cstore_fdw doesn't yet make optimizations related to query planning and execution. We made experiments in that direction (https://news.ycombinator.com/item?id=8423825), but making those changes production ready is no small effort.
Since all benchmarks in this blog post are for in-memory data, I don't know how much they would benefit from cstore. If I have the time, I'll give it a try and update this comment with the results.
Or am I missing the fact that these benchmarks are run on a reference spec which is comparatively old?
It has different isolation levels, some involving snapshots and some not. I think most concurrency issues are (by default) dealt with by locks, which start at row level and can escalate to page and table level (with significant slowdown seen when lock escalation happens in contentious places).
But that only has an effect if it's under write.
I managed an enterprise applications group that primarily used mssql for data in 2007. I can't recall which mvcc implementation they use. Our servers were MSSQL 2000 and I remember being a bit more than surprised when the DBAs told me the root of the performance problems we had were due to lock escalation; having come from the Oracle and PostgreSQL worlds, I was naive enough to have thought that lock escalation implementations like this were historical curiosities rather than something I'd actually run into... live and learn I guess.
So imagine a table with a single integer column and 1 billion rows. Instead of the 9 byte per row overhead (7 bytes for row metadata and 2 bytes for the page offset), you instead have 23 bytes, plus the 4 bytes to hold the int32. Without RCS on, that row would only have a 9 byte overhead and the rowsize would be 13 bytes.
so 1billion * 27bytes = 25.14 GB 1billion * 13bytes = 12.11 GB
There are some other performance tradeoffs (walking row version values, pagesplits when updating records without rowversions after converting the database, et cetera).
Ultimately, MVCC is great for contention, but it stinks if you're trying to efficiently pack in data.
For an alternative perspective, I sometimes bemoan the size of tables in postgres because of the mandatory overhead and versioning.
Oracle's MVCC method doesn't have that problem... but then you get the imfamous ORA-01555, "Snapshot too Old" from time to time.
Oh well, no perfect worlds I suppose.
Deleted comment
Deleted comment