andygrove commented on code in PR #2416:
URL:
https://github.com/apache/datafusion-ballista/pull/2416#discussion_r4116246953
##########
ballista/scheduler/src/state/mod.rs:
##########
@@ -305,23 +340,35 @@ impl<T: 'static + AsLogicalPlan, U: 'static +
AsExecutionPlan> SchedulerState<T,
.await
{
Ok(executor) => {
- if let Err(e) = state
+ match state
.task_manager
.launch_multi_task(&executor, tasks,
&state.executor_manager)
.await
{
- let err_msg = format!("Failed to launch new task:
{e}");
- error!("{}", err_msg.clone());
-
- // It's OK to remove executor aggressively,
- // since if the executor is in healthy state, it
will be registered again.
- state
- .remove_executor(&executor_id, Some(err_msg),
&sender)
- .await;
-
- false
- } else {
- true
+ Ok(unpreparable) => {
Review Comment:
This is its own PR now, #2492, and it has merged. Unpreparable tasks go
through the same path #2477 and #2016 use for rejected tasks, so their slots
are refunded and the job fails. With main merged in, the change no longer shows
up here.
##########
ballista/scheduler/src/state/task_manager.rs:
##########
@@ -1456,4 +1471,42 @@ mod tests {
assert_eq!(manager.running_job_number(), 0);
Ok(())
}
+
+ /// A task whose definition cannot be prepared is never sent to an
executor,
+ /// so no status will ever come back for it. `launch_multi_task` has to
name
+ /// the job it dropped, or the job stays Running for the life of the
+ /// scheduler and everything awaiting its status blocks with no timeout.
+ #[tokio::test]
+ async fn tasks_that_cannot_be_prepared_are_reported_against_their_job() ->
Result<()>
Review Comment:
This moved to #2492 along with the fix. The test there binds two jobs, makes
one job's tasks unpreparable, and checks that only that job is reported failed
and that exactly its slots are freed. It fails without the fix.
##########
ballista/core/src/serde/scheduler/mod.rs:
##########
@@ -71,11 +71,27 @@ impl PartitionLocation {
/// vocabulary of one writer. The stored field keeps that name until the
/// construction sites are swept; the wire protocol already speaks layouts.
pub fn layout(&self) -> ShuffleLayout {
- if self.is_sort_shuffle {
- ShuffleLayout::Sort
- } else {
- ShuffleLayout::Passthrough
- }
+ shuffle_layout(self.is_sort_shuffle)
+ }
+}
+
+impl crate::serde::protobuf::PartitionLocation {
Review Comment:
Done in #2495, which has merged. `fetch_partition` now uses `layout()` too,
so the rule lives in one place.
##########
ballista/flight-sql/src/service.rs:
##########
@@ -0,0 +1,852 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! The Arrow Flight SQL frontend itself.
+
+use std::pin::Pin;
+use std::sync::Arc;
+use std::time::Duration;
+
+use arrow::array::RecordBatch;
+use arrow::datatypes::{Schema, SchemaRef};
+use arrow::ipc::writer::IpcWriteOptions;
+use arrow_flight::encode::FlightDataEncoderBuilder;
+use arrow_flight::flight_service_server::FlightService;
+use arrow_flight::sql::server::FlightSqlService;
+use arrow_flight::sql::{
+ ActionCancelQueryRequest, ActionCancelQueryResult,
+ ActionClosePreparedStatementRequest, ActionCreatePreparedStatementRequest,
+ ActionCreatePreparedStatementResult, Any, CommandGetCatalogs,
CommandGetDbSchemas,
+ CommandGetSqlInfo, CommandGetTableTypes, CommandGetTables,
CommandGetXdbcTypeInfo,
+ CommandPreparedStatementQuery, CommandStatementQuery,
CommandStatementUpdate,
+ ProstMessageExt, SqlInfo, TicketStatementQuery,
+ metadata::{SqlInfoData, XdbcTypeInfoData},
+ server::PeekableFlightDataStream,
+};
+use arrow_flight::{
+ Action, FlightData, FlightDescriptor, FlightEndpoint, FlightInfo,
HandshakeRequest,
+ HandshakeResponse, IpcMessage, SchemaAsIpc, Ticket,
+};
+use ballista_core::error::BallistaError;
+use ballista_core::flight_proxy_service::BallistaFlightProxyService;
+use ballista_core::planner::scans_only_local_tables;
+use ballista_core::serde::protobuf::PartitionLocation;
+use ballista_core::serde::scheduler::{Action as BallistaAction,
ShuffleFileKind};
+use ballista_core::serde::{decode_protobuf, protobuf};
+use datafusion::logical_expr::{DdlStatement, LogicalPlan};
+use datafusion::prelude::SessionContext;
+use futures::{Stream, TryStreamExt};
+use prost::Message;
+use tonic::metadata::MetadataMap;
+use tonic::{Request, Response, Status, Streaming};
+use uuid::Uuid;
+
+use crate::auth::{AnonymousAuthenticator, Authenticator};
+use crate::backend::QueryBackend;
+use crate::metadata;
+use crate::session::{LocalResult, Prepared, SessionStore};
+use crate::ticket::StatementHandle;
+
+/// Session shared by every client that connects without authenticating.
+///
+/// Anonymous clients cannot be told apart, so they necessarily share catalog
+/// state. Configure an [`Authenticator`] to get a session per connection.
+pub const ANONYMOUS_SESSION: &str = "flight-sql-anonymous";
+
+/// How long a session, prepared statement, or unredeemed local result may sit
+/// idle before it is discarded.
+const DEFAULT_TTL: Duration = Duration::from_secs(30 * 60);
+
+/// How often expired handles are swept.
+const REAP_INTERVAL: Duration = Duration::from_secs(60);
+
+type DoGetStream =
+ Pin<Box<dyn Stream<Item = Result<FlightData, Status>> + Send + 'static>>;
+
+/// Serves Arrow Flight SQL on behalf of a Ballista cluster.
+///
+/// Clients send SQL text; the frontend plans it against the session's catalog,
+/// submits the plan through a [`QueryBackend`], and hands back one
+/// `FlightEndpoint` per output partition. `DoGet` on those tickets is proxied
+/// to the executor holding the partition, so clients never need to reach
+/// executors themselves — the failure mode that made the pre-46.0.0
+/// implementation unusable behind NAT, Docker, and Kubernetes.
+pub struct BallistaFlightSqlService<B: QueryBackend> {
+ backend: Arc<B>,
+ proxy: BallistaFlightProxyService,
+ auth: Arc<dyn Authenticator>,
+ store: Arc<SessionStore>,
+ sql_info: SqlInfoData,
+ xdbc_info: XdbcTypeInfoData,
+}
+
+impl<B: QueryBackend> BallistaFlightSqlService<B> {
+ /// Builds a frontend over `backend`, using `proxy` to stream partition
+ /// data back from executors.
+ ///
+ /// The service authenticates nobody until an [`Authenticator`] is supplied
+ /// via [`with_authenticator`](Self::with_authenticator).
+ pub fn new(backend: Arc<B>, proxy: BallistaFlightProxyService) -> Self {
+ let store = Arc::new(SessionStore::new(DEFAULT_TTL));
+
+ let backend_for_reaper = backend.clone();
+ store.spawn_reaper(REAP_INTERVAL, move |session_id| {
+ let backend = backend_for_reaper.clone();
+ async move {
+ if let Err(e) = backend.close_session(&session_id).await {
+ log::warn!("flight-sql: failed to close session
{session_id}: {e}");
+ }
+ }
+ });
+
+ Self {
+ backend,
+ proxy,
+ auth: Arc::new(AnonymousAuthenticator),
+ store,
+ sql_info: metadata::sql_info(),
+ xdbc_info: metadata::xdbc_type_info(),
+ }
+ }
+
+ /// Installs an authenticator. Without one, every handshake is accepted and
+ /// unauthenticated clients share a single session.
+ pub fn with_authenticator(mut self, auth: Arc<dyn Authenticator>) -> Self {
+ self.auth = auth;
+ self
+ }
+
+ /// True when the service will accept unauthenticated clients, which the
+ /// scheduler logs at startup.
+ pub fn allows_anonymous(&self) -> bool {
+ self.auth.allows_anonymous()
+ }
+
+ /// Resolves the Ballista session for a request from its bearer token.
+ fn session_id(&self, metadata: &MetadataMap) -> Result<String, Status> {
+ match bearer_token(metadata) {
+ Some(token) => self.store.session(&token).ok_or_else(|| {
+ Status::unauthenticated(
+ "unknown or expired session token; re-run the Flight
handshake",
+ )
+ }),
+ None if self.auth.allows_anonymous() =>
Ok(ANONYMOUS_SESSION.to_string()),
+ None => Err(Status::unauthenticated(
+ "missing bearer token; authenticate with the Flight handshake
first",
+ )),
+ }
+ }
+
+ /// Returns the context for `session_id`, building it on first use.
+ ///
+ /// The cache is what makes a session a session: [`QueryBackend::session`]
+ /// builds a fresh `SessionContext` every call, so without it a table
+ /// created by one request would be invisible to the next — and every
+ /// request would pay for a full DataFusion session to be constructed.
+ async fn open_session(
+ &self,
+ session_id: &str,
+ ) -> Result<Arc<SessionContext>, Status> {
+ if let Some(ctx) = self.store.context(session_id) {
+ return Ok(ctx);
+ }
+
+ let ctx = self
+ .backend
+ .session(session_id)
+ .await
+ .map_err(|e| Status::internal(format!("failed to open session:
{e}")))?;
+
+ Ok(self.store.insert_context(session_id.to_string(), ctx))
+ }
+
+ /// Resolves a request to the context it should be planned against.
+ async fn context(
+ &self,
+ metadata: &MetadataMap,
+ ) -> Result<Arc<SessionContext>, Status> {
+ self.open_session(&self.session_id(metadata)?).await
+ }
+
+ /// Plans `sql` against the session's catalog.
+ async fn plan(ctx: &SessionContext, sql: &str) -> Result<LogicalPlan,
Status> {
+ ctx.state()
+ .create_logical_plan(sql)
+ .await
+ .map_err(|e| Status::invalid_argument(format!("failed to plan
query: {e}")))
+ }
+
+ /// Runs a planned statement and describes where to collect its results.
+ async fn flight_info_for(
+ &self,
+ ctx: Arc<SessionContext>,
+ plan: LogicalPlan,
+ descriptor: FlightDescriptor,
+ job_name: &str,
+ ) -> Result<FlightInfo, Status> {
+ if let Disposition::Unsupported(reason) = disposition(&plan) {
+ return Err(Status::unimplemented(reason));
+ }
+
+ if disposition(&plan) == Disposition::RunOnScheduler {
+ let (schema, batches) = execute_locally(&ctx, plan).await?;
+ let handle = Uuid::new_v4().to_string();
+ self.store.insert_result(
+ handle.clone(),
+ LocalResult {
+ schema: schema.clone(),
+ batches,
+ },
+ );
+
+ let endpoint = Self::endpoint(StatementHandle::Local(handle));
+ return build_flight_info(&schema, vec![endpoint], descriptor);
+ }
+
+ let result = self
+ .backend
+ .execute(job_name, ctx, plan)
+ .await
+ .map_err(query_failed)?;
+
+ let endpoints = result
+ .partitions
+ .into_iter()
+ .map(|location| partition_handle(location).map(Self::endpoint))
+ .collect::<Result<Vec<_>, _>>()
+ .map_err(|e| Status::internal(format!("invalid partition location:
{e}")))?;
+
+ log::debug!(
+ "flight-sql: job {} produced {} endpoint(s)",
+ result.job_id,
+ endpoints.len()
+ );
+
+ build_flight_info(&result.schema, endpoints, descriptor)
+ }
+
+ /// Builds the endpoint a client redeems for one slice of the result.
+ ///
+ /// Endpoints carry no location, which Flight defines as "fetch from the
+ /// server that gave you this FlightInfo". That keeps every
cluster-internal
+ /// address off the wire and lets the frontend work unchanged behind NAT, a
+ /// load balancer, or an ingress.
+ fn endpoint(handle: StatementHandle) -> FlightEndpoint {
+ let ticket = TicketStatementQuery {
+ statement_handle: handle.encode().into(),
+ };
+ FlightEndpoint::new().with_ticket(Ticket {
+ ticket: ticket.as_any().encode_to_vec().into(),
+ })
+ }
+
+ /// Serves a metadata command by round-tripping the command itself as the
+ /// ticket, so `DoGet` lands back on the matching handler.
+ fn metadata_info<C: ProstMessageExt>(
+ command: C,
+ schema: SchemaRef,
+ descriptor: FlightDescriptor,
+ ) -> Result<Response<FlightInfo>, Status> {
+ let endpoint = FlightEndpoint::new().with_ticket(Ticket {
+ ticket: command.as_any().encode_to_vec().into(),
+ });
+ build_flight_info(&schema, vec![endpoint],
descriptor).map(Response::new)
+ }
+}
+
+/// What the frontend should do with a planned statement.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+enum Disposition {
+ /// Execute on the scheduler: it mutates session state and produces no data
+ /// worth distributing.
+ RunOnScheduler,
+ /// Submit to the cluster.
+ Distribute,
+ /// Refuse, with an explanation for the client.
+ Unsupported(&'static str),
+}
+
+/// Classifies a plan.
+///
+/// One function rather than two predicates, because the interesting cases are
+/// the ones where the answers overlap: `CREATE TABLE AS SELECT` is DDL, and
+/// DDL runs on the scheduler, so a caller that asked "is this DDL?" before
+/// asking "is this supported?" would silently execute its query on one node.
+/// Returning a single verdict makes that ordering impossible to get wrong.
+fn disposition(plan: &LogicalPlan) -> Disposition {
+ match plan {
+ LogicalPlan::Dml(_) => Disposition::Unsupported(
+ "Ballista Flight SQL does not support INSERT/UPDATE/DELETE; \
+ the distributed write path is not implemented",
+ ),
+ LogicalPlan::Copy(_) => Disposition::Unsupported(
+ "Ballista Flight SQL does not support COPY; \
+ the distributed write path is not implemented",
+ ),
+ LogicalPlan::Ddl(DdlStatement::CreateMemoryTable(_)) =>
Disposition::Unsupported(
+ "Ballista Flight SQL does not support CREATE TABLE AS SELECT, \
+ because it would execute on the scheduler rather than the
cluster; \
+ use CREATE EXTERNAL TABLE over data the executors can read",
+ ),
+ // Other DDL only edits the session catalog, and `SET`-style statements
+ // only edit session config; neither has anything to distribute.
+ LogicalPlan::Ddl(_) | LogicalPlan::Statement(_) =>
Disposition::RunOnScheduler,
+ // `SHOW ...` and anything else reading `information_schema` describes
+ // the catalog held here. There is nothing to distribute, and the
+ // physical form of those scans cannot be serialized for an executor,
+ // so distributing one would hand the client a job that never runs.
+ _ if scans_only_local_tables(plan) => Disposition::RunOnScheduler,
+ _ => Disposition::Distribute,
+ }
+}
+
+async fn execute_locally(
+ ctx: &SessionContext,
+ plan: LogicalPlan,
+) -> Result<(SchemaRef, Vec<RecordBatch>), Status> {
+ let df = ctx
+ .execute_logical_plan(plan)
+ .await
+ .map_err(|e| Status::internal(format!("failed to execute statement:
{e}")))?;
+ let planned_schema: SchemaRef = Arc::new(df.schema().as_arrow().clone());
+ let batches = df
+ .collect()
+ .await
+ .map_err(|e| Status::internal(format!("failed to execute statement:
{e}")))?;
+
+ // Prefer the schema the data actually carries; DataFusion's DDL results
+ // are empty and their frame schema is not always the same object.
+ let schema = batches
+ .first()
+ .map(|batch| batch.schema())
+ .unwrap_or(planned_schema);
+
+ Ok((schema, batches))
+}
+
+/// Turns a shuffle partition into the ticket payload the Flight proxy already
+/// knows how to redeem.
+fn partition_handle(
+ location: PartitionLocation,
+) -> Result<StatementHandle, BallistaError> {
+ let layout = location.layout();
+ let partition_id = location.partition_id.ok_or_else(|| {
+ BallistaError::Internal("partition location has no partition
id".to_string())
+ })?;
+ let executor = location.executor_meta.ok_or_else(|| {
+ BallistaError::Internal("partition location has no executor
metadata".to_string())
+ })?;
+
+ let action = BallistaAction::FetchPartition {
+ job_id: partition_id.job_id.into(),
+ stage_id: partition_id.stage_id as usize,
+ partition_id: partition_id.partition_id as usize,
+ host: executor.host,
+ port: executor.port as u16,
+ file_id: location.file_id,
+ layout,
+ // A Flight client wants the whole partition: the data file, in full.
+ file_kind: ShuffleFileKind::Data,
+ byte_ranges: vec![],
+ };
+
+ let encoded: protobuf::Action = action.try_into()?;
+ Ok(StatementHandle::Partition(encoded.encode_to_vec()))
+}
+
+fn build_flight_info(
+ schema: &Schema,
+ endpoints: Vec<FlightEndpoint>,
+ descriptor: FlightDescriptor,
+) -> Result<FlightInfo, Status> {
+ FlightInfo::new()
+ .try_with_schema(schema)
+ .map_err(|e| Status::internal(format!("failed to encode result schema:
{e}")))
+ .map(|info| info.with_descriptor(descriptor).with_endpoints(endpoints))
+}
+
+/// Streams a single metadata batch back to the client.
+fn one_batch_response(batch: RecordBatch) -> Response<DoGetStream> {
+ batch_response(batch.schema(), vec![batch])
+}
+
+fn batch_response(schema: SchemaRef, batches: Vec<RecordBatch>) ->
Response<DoGetStream> {
+ let stream = FlightDataEncoderBuilder::new()
+ .with_schema(schema)
+ .build(futures::stream::iter(batches.into_iter().map(Ok)))
+ .map_err(|e| Status::internal(format!("failed to encode results:
{e}")));
+
+ Response::new(Box::pin(stream) as DoGetStream)
+}
+
+fn bearer_token(metadata: &MetadataMap) -> Option<String> {
+ let value = metadata.get("authorization")?.to_str().ok()?;
+ value
+ .strip_prefix("Bearer ")
+ .or_else(|| value.strip_prefix("bearer "))
+ .map(str::to_string)
+}
+
+fn invalid_ticket(e: impl std::fmt::Display) -> Status {
+ Status::invalid_argument(format!("invalid ticket: {e}"))
+}
+
+fn query_failed(e: BallistaError) -> Status {
+ // A failed job is the client's problem to see, not an opaque 500.
+ Status::internal(format!("query execution failed: {e}"))
+}
+
+fn encode_schema(schema: &Schema) -> Result<Vec<u8>, Status> {
+ let message: IpcMessage = SchemaAsIpc::new(schema,
&IpcWriteOptions::default())
+ .try_into()
+ .map_err(|e| Status::internal(format!("failed to encode schema:
{e}")))?;
+ Ok(message.0.to_vec())
+}
+
+#[tonic::async_trait]
+impl<B: QueryBackend> FlightSqlService for BallistaFlightSqlService<B> {
+ type FlightService = BallistaFlightSqlService<B>;
+
+ async fn do_handshake(
+ &self,
+ request: Request<Streaming<HandshakeRequest>>,
+ ) -> Result<
+ Response<Pin<Box<dyn Stream<Item = Result<HandshakeResponse, Status>>
+ Send>>>,
+ Status,
+ > {
+ let identity = self.auth.authenticate(request.metadata()).await?;
+
+ let token = Uuid::new_v4().to_string();
+ let session_id = format!("flight-sql-{}", Uuid::new_v4());
+
+ // Build the session eagerly so a failure surfaces at handshake time
+ // rather than on the client's first query, and so the first query does
+ // not pay for it.
+ self.open_session(&session_id).await?;
+ self.store.insert_session(token.clone(), session_id.clone());
+
+ log::debug!(
+ "flight-sql: handshake for {:?} bound to session {session_id}",
+ identity.user
+ );
+
+ let result = HandshakeResponse {
+ protocol_version: 0,
+ payload: token.clone().into(),
+ };
+ let stream = futures::stream::once(async move { Ok(result) });
+
+ let mut response = Response::new(Box::pin(stream) as _);
+ response.metadata_mut().insert(
+ "authorization",
+ format!("Bearer {token}")
+ .parse()
+ .map_err(|_| Status::internal("failed to encode session
token"))?,
+ );
+ Ok(response)
+ }
+
+ /// Redeems tickets that are not Flight SQL commands.
+ ///
+ /// Ballista's own Rust client fetches shuffle output with a
+ /// `ballista.protobuf.Action` ticket. When the Flight SQL frontend is
+ /// mounted it replaces the standalone proxy on the scheduler's port, so it
+ /// has to keep serving those tickets.
+ async fn do_get_fallback(
Review Comment:
Done. The proxy now maps an undecodable ticket to `InvalidArgument` itself,
and `do_get_fallback` just hands the ticket to it. The proxy half of that
landed in #2495.
##########
ballista/core/src/planner.rs:
##########
@@ -165,6 +162,19 @@ impl<T: 'static + AsLogicalPlan> QueryPlanner for
BallistaQueryPlanner<T> {
}
}
+/// Returns `true` if every table `plan` scans lives in `information_schema`.
+///
+/// Such a plan reads the catalog of the node it is planned on and has nothing
+/// to distribute; worse, the physical form of those scans
(`StreamingTableExec`)
+/// cannot be serialized for an executor, so distributing one produces a job
+/// that can never run. Callers that hand plans to the cluster should answer
+/// these locally instead.
+pub fn scans_only_local_tables(plan: &LogicalPlan) -> bool {
+ let mut local_run = LocalRun::default();
+ let _ = plan.visit(&mut local_run);
Review Comment:
Good catch. I split this out as #2493, which switches `LocalRun` to
`visit_with_subqueries` and adds a test with this query. It's still open, so
the change stays in this diff until it merges.
--
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]