Hi +1 with Yufei and Robert. Thanks Yong for agreeing to run the benchmarks first.
Having concrete data from polaris-tools/benchmarks is a good before introducing multi-datasource and cache-isolation complexity in Polaris. It will show us whether connection exhaustion is a real bottleneck compared to metadata file IO, and whether pool multiplexers (like PgBouncer or RDS Proxy) might be good enough. Happy to help with the testing if needed. Thanks, Regards JB On Wed, Sep 30, 2026 at 1:18 AM Yong Zheng <[email protected]> wrote: > > Hello Yufei, > > I think that is a fair point. Let me get some setup and benchmarks up in the > upcoming weeks then I can update this thread again. > > Thanks, > Yong > > > > On Sep 29, 2026, at 12:38 PM, Yufei Gu <[email protected]> wrote: > > > > 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 > >>>>>>>>> > >>>>>>>> > >>>>>>> > >>>>>> > >>>>> > >>>> > >>> > >>
