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