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]
