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