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 >>>>>>>>> >>>>>>>> >>>>>>> >>>>>> >>>>> >>>> >>> >>
