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

Reply via email to