-
Notifications
You must be signed in to change notification settings - Fork 208
feat(java/driver/flight-sql): implement Flight SQL session management #4444
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
2ab7806
b69cb83
ee9b890
52eb8a2
cfcf3f5
7580eae
160cef4
f2341e4
72bfccd
d9f766f
b54f893
0d475c8
d06a170
4963d8c
55025e4
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -38,13 +38,22 @@ | |
| import org.apache.arrow.adbc.core.AdbcStatement; | ||
| import org.apache.arrow.adbc.core.AdbcStatusCode; | ||
| import org.apache.arrow.adbc.core.BulkIngestMode; | ||
| import org.apache.arrow.adbc.core.TypedKey; | ||
| import org.apache.arrow.adbc.sql.SqlQuirks; | ||
| import org.apache.arrow.flight.CallOption; | ||
| import org.apache.arrow.flight.CloseSessionRequest; | ||
| import org.apache.arrow.flight.FlightCallHeaders; | ||
| import org.apache.arrow.flight.FlightClient; | ||
| import org.apache.arrow.flight.FlightEndpoint; | ||
| import org.apache.arrow.flight.FlightRuntimeException; | ||
| import org.apache.arrow.flight.FlightStatusCode; | ||
| import org.apache.arrow.flight.GetSessionOptionsRequest; | ||
| import org.apache.arrow.flight.HeaderCallOption; | ||
| import org.apache.arrow.flight.Location; | ||
| import org.apache.arrow.flight.SessionOptionValue; | ||
| import org.apache.arrow.flight.SessionOptionValueFactory; | ||
| import org.apache.arrow.flight.SetSessionOptionsRequest; | ||
| import org.apache.arrow.flight.SetSessionOptionsResult; | ||
| import org.apache.arrow.flight.Ticket; | ||
| import org.apache.arrow.flight.auth2.BasicAuthCredentialWriter; | ||
| import org.apache.arrow.flight.client.ClientCookieMiddleware; | ||
|
|
@@ -208,11 +217,125 @@ public void setAutoCommit(boolean enableAutoCommit) throws AdbcException { | |
| } | ||
| } | ||
|
|
||
| @Override | ||
| public <T> T getOption(TypedKey<T> key) throws AdbcException { | ||
| final String k = key.getKey(); | ||
|
|
||
| if (k.equals(FlightSqlConnectionProperties.SESSION_OPTIONS)) { | ||
| if (key.getType() != String.class) { | ||
| return AdbcConnection.super.getOption(key); | ||
| } | ||
| return key.cast(FlightSqlSessionUtil.toJson(fetchSessionOptionsOrEmpty())); | ||
| } | ||
|
|
||
| final String prefix; | ||
| if (k.startsWith(FlightSqlConnectionProperties.SESSION_OPTION_BOOL_PREFIX)) { | ||
| prefix = FlightSqlConnectionProperties.SESSION_OPTION_BOOL_PREFIX; | ||
| } else if (k.startsWith(FlightSqlConnectionProperties.SESSION_OPTION_STRING_LIST_PREFIX)) { | ||
| prefix = FlightSqlConnectionProperties.SESSION_OPTION_STRING_LIST_PREFIX; | ||
| } else if (k.startsWith(FlightSqlConnectionProperties.SESSION_OPTION_PREFIX)) { | ||
| prefix = FlightSqlConnectionProperties.SESSION_OPTION_PREFIX; | ||
| } else { | ||
| return AdbcConnection.super.getOption(key); | ||
| } | ||
|
|
||
| final String name = k.substring(prefix.length()); | ||
| if (name.isEmpty()) { | ||
| throw AdbcException.invalidArgument("[Flight SQL] Session option name must not be empty"); | ||
| } | ||
| final Object raw = | ||
| FlightSqlSessionUtil.require(fetchSessionOptionsOrEmpty(), name) | ||
| .acceptVisitor(FlightSqlSessionUtil.TO_JAVA); | ||
| if (raw == null) { | ||
| throw new AdbcException( | ||
| "[Flight SQL] Session option not found: " + name, | ||
| null, | ||
| AdbcStatusCode.NOT_FOUND, | ||
| null, | ||
| 0); | ||
| } | ||
| final T result = FlightSqlSessionUtil.cast(key, raw, name); | ||
| return result != null ? result : AdbcConnection.super.getOption(key); | ||
| } | ||
|
|
||
| @Override | ||
| public <T> void setOption(TypedKey<T> key, T value) throws AdbcException { | ||
| final String k = key.getKey(); | ||
|
|
||
| if (k.startsWith(FlightSqlConnectionProperties.SESSION_OPTION_ERASE_PREFIX)) { | ||
| final String name = | ||
| k.substring(FlightSqlConnectionProperties.SESSION_OPTION_ERASE_PREFIX.length()); | ||
| doSetSessionOption(name, SessionOptionValueFactory.makeEmptySessionOptionValue()); | ||
|
|
||
| } else if (k.startsWith(FlightSqlConnectionProperties.SESSION_OPTION_BOOL_PREFIX)) { | ||
| if (value == null) { | ||
| throw invalidNullValue(k); | ||
| } | ||
| final String name = | ||
| k.substring(FlightSqlConnectionProperties.SESSION_OPTION_BOOL_PREFIX.length()); | ||
| final boolean b; | ||
| if (value instanceof Boolean) { | ||
| b = (Boolean) value; | ||
| } else { | ||
| b = FlightSqlSessionUtil.parseStrictBoolean(value.toString(), name); | ||
| } | ||
| doSetSessionOption(name, SessionOptionValueFactory.makeSessionOptionValue(b)); | ||
|
|
||
| } else if (k.startsWith(FlightSqlConnectionProperties.SESSION_OPTION_STRING_LIST_PREFIX)) { | ||
| if (value == null) { | ||
| throw invalidNullValue(k); | ||
| } | ||
| final String name = | ||
| k.substring(FlightSqlConnectionProperties.SESSION_OPTION_STRING_LIST_PREFIX.length()); | ||
| final String[] arr; | ||
| if (value instanceof String[]) { | ||
| arr = (String[]) value; | ||
| } else { | ||
| arr = FlightSqlSessionUtil.parseJsonArray(value.toString()); | ||
| } | ||
| doSetSessionOption(name, SessionOptionValueFactory.makeSessionOptionValue(arr)); | ||
|
|
||
| } else if (k.startsWith(FlightSqlConnectionProperties.SESSION_OPTION_PREFIX)) { | ||
| if (value == null) { | ||
| throw invalidNullValue(k); | ||
| } | ||
| final String name = k.substring(FlightSqlConnectionProperties.SESSION_OPTION_PREFIX.length()); | ||
| final SessionOptionValue sv; | ||
| if (value instanceof Long) { | ||
| sv = SessionOptionValueFactory.makeSessionOptionValue((Long) value); | ||
| } else if (value instanceof Double) { | ||
| sv = SessionOptionValueFactory.makeSessionOptionValue((Double) value); | ||
| } else { | ||
| sv = SessionOptionValueFactory.makeSessionOptionValue(value.toString()); | ||
| } | ||
| doSetSessionOption(name, sv); | ||
|
|
||
| } else if (k.equals(FlightSqlConnectionProperties.SESSION_OPTIONS)) { | ||
| throw AdbcException.notImplemented( | ||
| "[Flight SQL] adbc.flight.sql.session.options is read-only"); | ||
|
|
||
| } else { | ||
| AdbcConnection.super.setOption(key, value); | ||
| } | ||
| } | ||
|
|
||
| private static AdbcException invalidNullValue(String key) { | ||
| return AdbcException.invalidArgument( | ||
| "[Flight SQL] null value not allowed for key: " | ||
| + key | ||
| + " - use adbc.flight.sql.session.optionerase.<name> to erase an option"); | ||
| } | ||
|
|
||
| @Override | ||
| public void close() throws AdbcException { | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I think this could be cleaner if you just passed all of these into AutoCloseables.close (AutoCloseable can be used as a FunctionalInterface)
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Kept closeSession out of the AutoCloseables.close call since it's an RPC that needs selective error-swallowing rather than being a plain closeable resource, but everything that actually is one (clientCache, client, allocator) now goes through a single flat call. |
||
| clientCache.invalidateAll(); | ||
| try { | ||
| AutoCloseables.close(client, allocator); | ||
| // Best-effort: the Go driver also ignores all errors closing the session. | ||
| client.closeSession(new CloseSessionRequest(), callOptions); | ||
| } catch (FlightRuntimeException e) { | ||
| // ignore | ||
| } | ||
| try { | ||
| AutoCloseables.close(clientCache::invalidateAll, client, allocator); | ||
| } catch (Exception e) { | ||
| throw AdbcException.internal("[Flight SQL] Failed to close connection").withCause(e); | ||
| } | ||
|
|
@@ -223,6 +346,39 @@ public String toString() { | |
| return "FlightSqlConnection{" + "client=" + client + '}'; | ||
| } | ||
|
|
||
| private Map<String, SessionOptionValue> fetchSessionOptionsOrEmpty() throws AdbcException { | ||
| try { | ||
| return client.getSessionOptions(new GetSessionOptionsRequest()).getSessionOptions(); | ||
| } catch (FlightRuntimeException e) { | ||
| // Go also treats INVALID_ARGUMENT as "server doesn't support sessions" here. | ||
| if (e.status().code() == FlightStatusCode.UNIMPLEMENTED | ||
| || e.status().code() == FlightStatusCode.INVALID_ARGUMENT) { | ||
| return Collections.emptyMap(); | ||
| } | ||
| throw FlightSqlDriverUtil.fromFlightException(e); | ||
| } | ||
| } | ||
|
|
||
| private void doSetSessionOption(String name, SessionOptionValue value) throws AdbcException { | ||
| if (name.isEmpty()) { | ||
| throw AdbcException.invalidArgument("[Flight SQL] Session option name must not be empty"); | ||
| } | ||
| final SetSessionOptionsResult result; | ||
| try { | ||
| result = | ||
| client.setSessionOptions( | ||
| new SetSessionOptionsRequest(Collections.singletonMap(name, value))); | ||
| } catch (FlightRuntimeException e) { | ||
| throw FlightSqlDriverUtil.fromFlightException(e); | ||
| } | ||
| if (result.hasErrors()) { | ||
| final SetSessionOptionsResult.Error err = result.getErrors().get(name); | ||
| final String errType = (err != null) ? err.value.name() : "UNKNOWN"; | ||
| throw AdbcException.invalidArgument( | ||
| "[Flight SQL] Failed to set session option '" + name + "': " + errType); | ||
| } | ||
| } | ||
|
|
||
| /** | ||
| * Initialize cached data to share between connections and create, test, and authenticate the | ||
| * first connection. | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.