andygrove opened a new pull request, #2507:
URL: https://github.com/apache/datafusion-ballista/pull/2507

   # Which issue does this PR close?
   
   No separate issue. Follow-up to #2493, which was [approved as a 
stopgap](https://github.com/apache/datafusion-ballista/pull/2493#pullrequestreview-5332101088)
 until there was a better answer to the mismatch between the client's and the 
scheduler's `information_schema`.
   
   # Rationale for this change
   
   The client's catalog is the one that counts. The client plans every query 
and ships its tables by value. A `ListingTable` scan carries its paths, schema 
and format, and the scheduler rebuilds it without looking anything up, so the 
scheduler's catalog never matters for a query the client planned.
   
   `information_schema` is the exception. Its tables describe the catalog of 
the session that planned the query, so they cannot be shipped by value, and 
resolving them by name on the scheduler would read the scheduler's catalog, 
which is the wrong one. That is why #2493 refuses a query that reads 
`information_schema` together with other tables, such as:
   
   ```sql
   SELECT table_name, (SELECT count(*) FROM test) AS row_count
   FROM information_schema.tables WHERE table_name = 'test'
   ```
   
   This PR runs that query on the cluster instead. Before a plan leaves the 
session that planned it, the parts that read only `information_schema` run 
there, and their rows go into the plan as constants. The scheduler never 
resolves `information_schema` itself, so the two catalogs no longer need to 
agree.
   
   # What changes are included in this PR?
   
   - `inline_information_schema(plan, planner, session)`, a new public function 
in `ballista-core`. It finds each largest subtree that reads only 
`information_schema`, runs it with the given planner and session, and replaces 
it with a `Values` node under a `Projection` that restores the subtree's 
qualified column names. Filters and aggregates over `information_schema` run on 
the client, so only their results travel.
     - An empty result becomes one placeholder row under `LIMIT 0`, because 
neither an empty `Values` nor an `EmptyRelation` keeps its schema through 
serialization.
     - A result without columns keeps its row count through an empty projection.
     - Everything uses stock datafusion-proto, so there is no codec change.
   - A subtree is not inlined if it refers to an outer query, contains a 
subquery expression, or contains an extension, DML, DDL, `COPY`, statement, 
`EXPLAIN` or `ANALYZE` node. For a correlated subquery, only the scan under the 
correlated filter is inlined.
   - `BallistaQueryPlanner` still runs plans that read only 
`information_schema` on the client. Any other plan goes through 
`inline_information_schema` before it is sent to the scheduler, `EXPLAIN` and 
`EXPLAIN ANALYZE` included.
   - `scans_only_information_schema` returns `bool`. Its only error was the 
refusal.
   - Unit tests compare each rewritten plan with DataFusion's answer for the 
original, before and after a round trip through 
`BallistaLogicalExtensionCodec`. They cover subqueries, correlated subqueries, 
joined `information_schema` tables, views, empty results, results without 
columns and repeated parts, and they plan with 
`SessionConfig::new_with_ballista()` so the plans have the shape 
`BallistaQueryPlanner` receives.
   - End-to-end tests, standalone and remote, for the query above, a join with 
`information_schema.columns`, a query that matches no rows, and 
`information_schema.df_settings` showing the client's settings. The refusal 
test from #2493 is removed.
   
   # Are there any user-facing changes?
   
   Queries that read `information_schema` together with other tables now run on 
the cluster and see the client's catalog, instead of failing. `EXPLAIN` of such 
a query shows a `Values` node where the `information_schema` part was, because 
that is what the cluster runs.
   
   The inlined rows travel inside the plan, so a query that joins a very large 
part of `information_schema.columns` can hit the gRPC message size limit 
(`ballista.client.grpc_max_message_size`, 16 MiB by default).
   
   Adds `inline_information_schema` to `ballista-core` and changes the return 
type of `scans_only_information_schema` from `Result<bool>` to `bool`. Neither 
has been in a release.
   
   The Flight SQL frontend in #2416 needs the same call before it distributes a 
plan. Whichever of the two PRs merges second will add it.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to