Can we take a step back and clarify the overall workload expectations first? - Expected Polaris request rate and QPS/TPS - Number of tables per namespace and catalog - Latency requirements
Once we have those baseline metrics, we can measure connection utilization and wait times, metastore latency, and metadata-file I/O to identify the actual bottleneck. As I mentioned in my previous email, For example, loadTable also reads metadata files, and that I/O could take much longer than getting a connection from the pool. The constraint of single Postgres instance connection count may or may not be the system bottleneck. Yufei On Tue, Sep 29, 2026 at 3:24 AM Robert Stupp <[email protected]> wrote: > Hi all, > > The workload described in the Medium post used 15 Polaris pods, reaching > ~400 sustained requests/second. > That is <27 requests/second per pod. It also appears to have mostly idle > connection pools and a small, static catalog. This does appear to prove a > JDBC connection saturation. > > Before adding a second datasource path, we should understand the > consistency behavior exposed to clients when primary and read-replica > states differ. > In particular, a replica can lag after a successful mutation, and the > EntityCache/Resolver path needs a defined policy for which state it may > use. > > An HTTP request header also cannot by itself establish that the full > operation is safe to serve from a replica. > Some apparently read-only service requests may have explicit or implicit > persistence side effects, or may need current primary state. > That needs a narrow, reviewed routing rule rather than relying on the HTTP > method or a header alone. > > I think we should first measure a representative, reproducible workload, > ideally using the Polaris Benchmarks tool [1]. > That should include realistic catalog size and request mix, pool > utilization and wait time, database time, metadata I/O, and replica lag. > Then we can define a narrow stale-read contract with isolated cache > semantics. That may validate a deliberately limited replica-serving design, > which could lead to a PR. > > It would also be useful to compare this with the existing NoSQL cache > approach. > Its cross-pod invalidation lets valid hot entity data be served without a > backend read. > That is not a substitute for a defined current-state contract, but it may > show whether cacheable read load, rather than JDBC connection capacity, is > the actual problem to solve. > > More generally, this is a useful test for the persistence contract: > Caller-visible operation semantics and current-state guarantees should not > be inferred from HTTP transport details or from a single JDBC deployment > topology. > > Robert > > [1] https://github.com/apache/polaris-tools/tree/main/benchmarks > > On Tue, Sep 29, 2026 at 9:28 AM Yong Zheng <[email protected]> wrote: > > > Hello Prithvi, > > > > Thanks for the summary, and thanks everyone for the feedback and > > suggestions. > > > > I’m not sure whether companies typically publish enough details about > > their production benchmarks and workload characteristics to validate this > > with real-world data. But as a practical example, with AWS RDS using a DB > > cluster or even a Multi-AZ setup, the read replica can otherwise sit > mostly > > idle while still being required for production availability. Being able > to > > use that capacity for read-heavy workloads could be useful. > > > > I’ll check whether I can share some relevant workload information and get > > back to the group. > > > > In the meantime, should I proceed with a PR based on the direction we’ve > > discussed so far? My understanding of the plan is: > > 1. First, keep the existing primary-only behavior as the default. > > 2. Add an optional replica datasource. > > 3. Make the routing decision at the persistence-operation level rather > > than based purely on HTTP method. > > 4. Keep mutations, auth, and operations with write/consistency > > requirements on the primary. > > 5. Allow explicitly opted-in read workloads to use the replica and > > document the potential replica lag (out-of-band). > > 6. Do not silently fall back to the primary if a replica operation cannot > > be served. > > 7. Avoid using the shared entity cache on the replica path. > > > > Please let me know if this matches the direction you have in mind, or if > > we should clarify anything before I start the PR. > > > > Thanks, > > Yong Zheng > > > > On 2026/09/24 20:01:38 Prithvi S wrote: > > > Hi all, > > > > > > Yong, I think we are on the same page about the retry case. If a read > > hits > > > the replica before the commit is visible, and the client assumes the > > write > > > failed and retries, you can ingest twice. Routing from the HTTP verb > > would > > > do that, and I would avoid it. Your original split into an ingestion > > > workload and a query workload still sounds right to me. Ingestion talks > > to > > > a Polaris on the primary. The query fleet talks to a Polaris that may > > read > > > from the replica. The client picks that by catalog URI, which Spark, > > Flink, > > > Trino, and PyIceberg already have. Auth stays on the primary, and so > does > > > anything that writes. listNamespaces and loadTable raise events, and > with > > > the JDBC persistence listener enabled those flush as inserts on their > own > > > session, so the serving side still needs the primary for that write. > JB's > > > suggestion of making the choice beside the persistence call is where I > > can > > > see it working, because that is where we know whether the call writes. > > If a > > > write does reach the replica, I would fail it in Polaris with a clear > > 4xx, > > > and not fall back to the primary or return a Postgres read-only error. > > > > > > Dmitri, I agree with you that a header is reasonable when the caller > > knows > > > they can live with lag. Polaris-Out-Of-Band is a clearer name than > > > Polaris-Readonly. Readonly can sound like the call is simply safe, and > > > out-of-band says the caller is accepting a stale read. It is the same > > kind > > > of opt-in as pointing the query fleet at the reader endpoint, for > someone > > > who can set it per request. The part I would be careful with is the one > > you > > > already flagged: not every GET can tolerate stale data. loadTable right > > > after commitTable is a GET. So I would treat the header as "this caller > > > accepts lag", and still decide at the persistence method whether that > > call > > > can run on the replica. Event flushes would stay on the primary. If the > > > header is set on a call we cannot serve that way, including a POST, I > > would > > > rather return a 4xx than ignore the header and run it on the primary, > so > > it > > > is clear the switch did not apply. For Spark and the other standard > > > clients, the catalog URI is probably still the practical switch, since > > they > > > will not know to send this header on some calls and not others. > > > > > > On the cache, your concern makes sense to me. > JdbcMetaStoreManagerFactory > > > keeps one InMemoryEntityCache per realm and shares it across requests. > An > > > entry stays for an hour after last use, and we keep the higher entity > > > version. A replica read can leave an old version in that cache, and a > > later > > > primary read can serve it. The other direction is awkward too: a > primary > > > read fills the cache, and an out-of-band read gets the fresh row and > does > > > not see the lag the caller asked for. Separate caches per pool would > keep > > > those apart. On the replica path I would lean toward turning the cache > > off, > > > which is the other option you mentioned, because a cache there can > hold a > > > stale row longer than the replica itself is behind. The primary cache > can > > > stay as it is. > > > > > > Yufei, I think your question should come first. Before adding a second > > > datasource, it would help to see a workload like the one Yong > described, > > > with many tables and frequent maintenance checks. The Anand numbers I > > cited > > > were one run, with an idle pool and a lot of listNamespaces on a quiet > > > catalog. On a maintenance scan I would look at connections in use > against > > > quarkus.datasource.jdbc.max-size, time waiting for a connection, time > in > > > the database, and time in metadata-file I/O inside loadTable. If > requests > > > are hardly waiting on a connection and loadTable is mostly file I/O, I > am > > > not sure a replica buys much, and PgBouncer or RDS Proxy may be the > > smaller > > > step for the session count. If the primary connections are what is > > actually > > > full, then an optional second datasource seems worth writing up: off > > unless > > > configured, used by the serving deployment, split at the persistence > > > method, with the lag documented and no entity cache on the replica > path. > > > > > > WDYT? > > > > > > Cheers, > > > Prithvi S > > > > > > On Mon, Sep 14, 2026 at 10:30 PM Yufei Gu <[email protected]> > wrote: > > > > > > > I think consistency needs to come first. Clients should be able to > read > > > > their changes after a successful write, regardless of how we route > > > > requests. > > > > > > > > Have we actually seen database connection limits become a bottleneck > > for > > > > reads? The connection-count calculation suggests a possible limit, > but > > it > > > > would help to check what slows us down first. For example, loadTable > > also > > > > reads metadata files, and that I/O could take much longer than > getting > > a > > > > connection from the pool. We might hit that bottleneck before running > > out > > > > of connections, in which case read replicas may not help much. > > > > > > > > Could we test this with a realistic workload and look at how many > > > > connections are in use, how long requests wait for a connection, and > > how > > > > much time goes into database queries versus metadata-file I/O? That > > would > > > > give us a clearer idea of whether read replicas would help and if the > > added > > > > complexity is justified. > > > > > > > > Thanks, > > > > Yufei > > > > > > > > On Mon, Sep 14, 2026 at 7:59 AM Dmitri Bourlatchkov < > [email protected]> > > > > wrote: > > > > > > > > > Hi Yong, Prithvi, > > > > > > > > > > I think Yong's header idea is reasonable for some deployments. I > > agree > > > > that > > > > > in general routing R/O requests to a replica can > > > > > violate clients' consistency expectations. However, if in a > > particular > > > > > situation the client (admin user) is aware of the expectations, > using > > > > this > > > > > as an optimization switch can be acceptable. > > > > > > > > > > Also, all GET requests are read-only by definition, but not all > GETs > > may > > > > be > > > > > able to tolerate reading stale data. > > > > > > > > > > Perhaps the matter can be clarified by using a different header > name > > to > > > > > highlight the implications. Something like "Polaris-Out-Of-Band: > > true". > > > > In > > > > > this form the header is applicable to all requests, but will alter > > > > > execution only for GETs. > > > > > > > > > > WDYT? > > > > > > > > > > Implementation-wise, if we have two different connection pools, > > > > > the InMemoryEntityCache may become more of an obstacle than an aid > > since > > > > > replication delays can cause odd effects in the cache (if it is > > shared > > > > > between the two connection pools). For the sake of sanity, it may > be > > > > > necessary to use different entity caches for each connection pool > > (or not > > > > > use the cache at all [1]). > > > > > > > > > > [1] > https://lists.apache.org/thread/tl6z48gdblsko3x0b9nt2917wb2fgtor > > > > > > > > > > Cheers, > > > > > Dmitri. > > > > > > > > > > On Sun, Sep 13, 2026 at 11:32 PM Yong Zheng <[email protected]> > > wrote: > > > > > > > > > > > Hello, > > > > > > > > > > > > Thanks for the quick review Prithvi. Yes around the rps as those > > are > > > > only > > > > > > provided as a quick math to show the scaling issue that people > can > > face > > > > > (or > > > > > > may already faced) when having a large lakehouse in an enterprise > > > > > > environment. > > > > > > > > > > > > I do agreed that we will make the RO DB endpoint optional (so > > existed > > > > > > deployment will continue to behavior the same way) and can be add > > when > > > > > > there is a needed for scaling by splitting readonly requests into > > this > > > > > > optional RO DB endpoint. However, I don't think this is a good > > idea to > > > > > make > > > > > > the server blindly route non-GET requests to default DB endpoint > > and > > > > GET > > > > > > requests to optional RO DB endpoint. As called out in your > > response, > > > > > > replication can have latency and it can have big impacts when > those > > > > > > happened and caused un-intensional retry from custom Iceberg > client > > > > (e.g. > > > > > > if people are doing a write then checking if write completed, > then > > > > retry > > > > > if > > > > > > write was not committed...in this case, latency can cause > > un-necessary > > > > > > retry as well as potentially dup ingestion for custom Iceberg > > writer). > > > > > > > > > > > > As we do use custom header to route requests to different REALM > in > > > > > > Polaris, I think it is safer for client to decide if they should > > route > > > > > > requests to RO endpoint via another customer header as oppose as > > let > > > > > server > > > > > > handles this blindly as described above. The workload used by > > Anand's > > > > is > > > > > > not very likely a real world workload or may only covered part of > > the > > > > use > > > > > > cases on how people are using Iceberg. I don' think we should use > > that > > > > > as a > > > > > > way to say 60% of the workload is listNamespaces which doesn't > > change. > > > > > For > > > > > > deployment where there are high number of tables (e.g. tenant > level > > > > > tables > > > > > > under different namespaces) and custom data-ops, checking if > table > > > > > > maintenance among all tables is critical and those do does very > > > > frequent. > > > > > > Thoughts? > > > > > > > > > > > > Thanks, > > > > > > Yong Zheng > > > > > > > > > > > > > > > > > > > > > > > > On 2026/09/13 11:44:27 Prithvi S wrote: > > > > > > > Hi Yong, > > > > > > > > > > > > > > Thanks for writing this up, and for connecting it with Anand's > > > > > load-test > > > > > > > write-up. The connection math is a good place to start. With > JDBC > > > > > > > persistence, the Agroal pool is the limit on each pod > > > > > > > (quarkus.datasource.jdbc.max-size, commonly 50), and > > max_connections > > > > on > > > > > > the > > > > > > > primary is the limit on the fleet. That is a real scaling axis, > > and > > > > > using > > > > > > > read replicas is worth discussing. > > > > > > > > > > > > > > I would be careful about treating ~5,333 rps as a planning > > target :) > > > > > > > Anand's run was ~400 rps sustained and ~800 rps peak on 15 > pods, > > with > > > > > low > > > > > > > CPU and well under one in-flight request per pod on average. > The > > > > > > > 50-connection pool was not what was biting, most of those > > connections > > > > > > were > > > > > > > idle. Scaling from there is likely to hit database CPU, I/O, > > WAL, or > > > > > lock > > > > > > > contention before a clean 5K-connection wall. There are also > > cheaper > > > > > > levers > > > > > > > in front of replicas.. right-sizing pool and pod count, and > > > > > multiplexing > > > > > > > with something like PgBouncer or RDS Proxy so pod count and > > database > > > > > > > sessions do not grow 1:1. > > > > > > > > > > > > > > On Polaris-Readonly, I am a bit hesitant to make this a > > client-side > > > > > > routing > > > > > > > decision. main concern is consistency. Iceberg clients can > > expect a > > > > > read > > > > > > > after a write in the same session to see that write, for > example > > > > > > > commitTable followed by loadTable. RDS/Aurora replicas are > > > > > asynchronous, > > > > > > so > > > > > > > sending those reads to a replica is a change in consistency > > model. > > > > That > > > > > > > should be an explicit contract, not something controlled by a > > header. > > > > > > > > > > > > > > There is also a practical issue with standard Iceberg clients. > > Spark, > > > > > > > Flink, Trino, and PyIceberg do not know about > Polaris-Readonly. A > > > > > static > > > > > > > header applies to the whole catalog, so a stray mutation, or > > > > something > > > > > > like > > > > > > > an Iceberg metrics report, could land on the replica too. And > > some > > > > > > > operations that look like reads can still write in Polaris > > (persisted > > > > > > > events, idempotency bookkeeping). The header does not turn > those > > side > > > > > > > effects off. If a mutation does hit the replica, I agree we > > should > > > > not > > > > > > > silently fall back to the primary, but Polaris should fail that > > > > itself > > > > > > with > > > > > > > a clear 4xx, not with a Postgres read-only error. > > > > > > > > > > > > > > I think a better path is to keep today's behavior as the > default > > and > > > > > make > > > > > > > replica use opt-in on the server: > > > > > > > > > > > > > > - Add an optional second Quarkus datasource pointing at the > > replica. > > > > > > > - Keep mutations, auth, and anything that must be linearizable > > with a > > > > > > > prior write on the primary; pure metadata reads may use the > > > > replica. > > > > > > > - Do not silently fall back to the primary if a replica call > > fails. > > > > > > > - Document that opted-in reads can see replica lag. > > > > > > > > > > > > > > That also fits the two-workload setup you described. The > > ingestion > > > > > fleet > > > > > > > can use the writer endpoint and the serving fleet the reader > > > > endpoint. > > > > > > The > > > > > > > mixed case (token issue on RW, catalog reads on RO) can then be > > > > handled > > > > > > > inside Polaris, rather than asking Iceberg clients to > understand > > a > > > > new > > > > > > > header. One other point from Anand's article, around 60% of the > > 30M > > > > > > > requests were listNamespaces against a catalog that barely > > changes. > > > > > > > Replicas would absorb that, but so would connector-side caching > > or a > > > > > > longer > > > > > > > refresh interval. > > > > > > > > > > > > > > I think this server-side dual-datasource approach is the right > > way to > > > > > > move > > > > > > > forward. WDYT? > > > > > > > > > > > > > > Regards, > > > > > > > Prithvi S > > > > > > > > > > > > > > On Sun, Sep 13, 2026 at 9:45 AM Yong Zheng <[email protected]> > > > > wrote: > > > > > > > > > > > > > > > Hello, > > > > > > > > > > > > > > > > When using JDBC as the backend metastore for Polaris, the > JDBC > > > > > > connection > > > > > > > > capacity of the active primary can become a limiting factor > > for the > > > > > > > > throughput of a given Polaris deployment. Taking RDS with > > > > PostgreSQL > > > > > > as an > > > > > > > > example, if we are using an AWS db.m9g.4xlarge (16 cores and > > 64 GB) > > > > > as > > > > > > the > > > > > > > > backend store, we can get up to 5K connections by default. > AWS > > uses > > > > > > > > "LEAST({DBInstanceClassMemory/9531392}, 5000)" for the > default > > > > > > PostgreSQL > > > > > > > > "max_connections", and this can be increased manually. > > > > > > > > > > > > > > > > Now, assuming we put 50 connections per pod, we can have up > to > > 100 > > > > > pods > > > > > > > > max (in reality, it will be less as a couple of connections > are > > > > > > reserved > > > > > > > > for the superuser, but let's stick with 100 pods to make the > > math > > > > > > simpler). > > > > > > > > > > > > > > > > With 50 connections per pod, this matches what Anand reported > > in > > > > > > > > > > > > > > > > > > > > > > > > > > https://medium.com/@obelix74/a-30-million-request-load-test-for-apache-polaris-fb5690040154 > > > > > > . > > > > > > > > If we assume linear scaling from the benchmark (which may not > > hold > > > > > due > > > > > > to > > > > > > > > database CPU, I/O, locking, and other bottlenecks), the > > theoretical > > > > > > > > throughput would be ~5,333 requests per second. This is > great, > > but > > > > if > > > > > > we > > > > > > > > ever want to handle more throughput, we would have no other > > option > > > > > > other > > > > > > > > than scaling up the RDS instance and overriding the default > max > > > > > allowed > > > > > > > > connections on the DB server when using the native AWS > > solution. > > > > > There > > > > > > are > > > > > > > > other solutions out there, such as using a multi-master > > deployment > > > > > for > > > > > > the > > > > > > > > backend DB or switching to a NoSQL backend, which we can > scale > > up a > > > > > lot > > > > > > > > easier (e.g. MongoDB). > > > > > > > > > > > > > > > > While there are solutions to scale up the infrastructure to > > support > > > > > > more > > > > > > > > connections, one particular thing that caught my eye is that > > we are > > > > > not > > > > > > > > using read-only replicas at all. > > > > > > > > > > > > > > > > Assuming we have two primary workloads: > > > > > > > > > > > > > > > > 1. Ingestion layer: this has both read and write. > > > > > > > > 2. Query/Serving layer: this is read-only queries to power > > various > > > > > > > > dashboards. > > > > > > > > > > > > > > > > We could also have two DB connection endpoints: > > > > > > > > > > > > > > > > 1. RW connection endpoint (active primary) > > > > > > > > 2. RO connection endpoint (read replicas) > > > > > > > > > > > > > > > > By default, everything should go to the RW endpoint. This > > matches > > > > the > > > > > > > > current workflow we support. However, if we know a > > query/serving > > > > > layer > > > > > > is > > > > > > > > read-only, we should route those requests to the RO > connection > > > > > > endpoint; > > > > > > > > the auth token request/renewal would still be fulfilled > > through the > > > > > RW > > > > > > > > connection endpoint for the query/serving layer. > > > > > > > > > > > > > > > > Now, to decide where a client should be sending requests to > > RW/RO, > > > > we > > > > > > can > > > > > > > > check the following: > > > > > > > > > > > > > > > > 1. Is this an auth request? If yes, always use the RW > endpoint. > > > > > > > > 2. Does this request have a special header (e.g. > > > > "Polaris-Readonly", > > > > > > which > > > > > > > > defaults to false or unset)? If yes, offload the request to > > the RO > > > > > > endpoint. > > > > > > > > > > > > > > > > In this case, if a non-read-only request is ever sent by the > > client > > > > > by > > > > > > > > accident, it will be failed by the backend DB because writes > > are > > > > not > > > > > > > > supported on RO replicas, and this is expected behavior. We > > should > > > > > not > > > > > > > > silently fall back to the RW endpoint in this case. > > > > > > > > > > > > > > > > With this approach, we can squeeze higher requests per second > > out > > > > of > > > > > a > > > > > > > > given setup by offloading the read-heavy query/serving > workload > > > > from > > > > > > the > > > > > > > > primary. If the load from the query/serving layer is > > significant, > > > > > this > > > > > > > > could give us another dimension to scale the infrastructure > > without > > > > > > having > > > > > > > > to keep scaling up the primary DB. > > > > > > > > > > > > > > > > Thanks, > > > > > > > > Yong Zheng > > > > > > > > > > > > > > > > > > > > > > > > > > > > > > > > > > > >
