Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions java/driver/flight-sql/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,11 @@
<artifactId>caffeine</artifactId>
<version>3.2.4</version>
</dependency>
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
<version>${dep.jackson.version}</version>
</dependency>
<dependency>
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java</artifactId>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,12 +26,18 @@
import org.apache.arrow.flight.CallOption;
import org.apache.arrow.flight.CancelFlightInfoRequest;
import org.apache.arrow.flight.CancelFlightInfoResult;
import org.apache.arrow.flight.CloseSessionRequest;
import org.apache.arrow.flight.CloseSessionResult;
import org.apache.arrow.flight.FlightDescriptor;
import org.apache.arrow.flight.FlightEndpoint;
import org.apache.arrow.flight.FlightInfo;
import org.apache.arrow.flight.FlightStream;
import org.apache.arrow.flight.GetSessionOptionsRequest;
import org.apache.arrow.flight.GetSessionOptionsResult;
import org.apache.arrow.flight.RenewFlightEndpointRequest;
import org.apache.arrow.flight.SchemaResult;
import org.apache.arrow.flight.SetSessionOptionsRequest;
import org.apache.arrow.flight.SetSessionOptionsResult;
import org.apache.arrow.flight.Ticket;
import org.apache.arrow.flight.sql.CancelResult;
import org.apache.arrow.flight.sql.FlightSqlClient;
Expand Down Expand Up @@ -287,6 +293,20 @@ public FlightEndpoint renewFlightEndpoint(
return client.renewFlightEndpoint(request, combine(options));
}

public SetSessionOptionsResult setSessionOptions(
SetSessionOptionsRequest request, CallOption... options) {
return client.setSessionOptions(request, combine(options));
}

public GetSessionOptionsResult getSessionOptions(
GetSessionOptionsRequest request, CallOption... options) {
return client.getSessionOptions(request, combine(options));
}

public CloseSessionResult closeSession(CloseSessionRequest request, CallOption... options) {
return client.closeSession(request, combine(options));
}

@Override
public void close() throws Exception {
AutoCloseables.close(client);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -96,7 +105,7 @@
(@Nullable Location key,
@Nullable FlightSqlClientWithCallOptions value,
RemovalCause cause) -> {
if (value == null || value == primaryClient) return;

Check warning on line 108 in java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlConnection.java

View workflow job for this annotation

GitHub Actions / Java 21/Linux with ErrorProne

[ReferenceEquality] Comparison using reference equality instead of value equality

Check warning on line 108 in java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlConnection.java

View workflow job for this annotation

GitHub Actions / Java 25/Linux with ErrorProne

[ReferenceEquality] Comparison using reference equality instead of value equality
try {
value.close();
} catch (Exception ex) {
Expand Down Expand Up @@ -208,11 +217,174 @@
}
}

@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");
}
if (!FlightSqlSessionUtil.supportsType(key, prefix)) {
return AdbcConnection.super.getOption(key);
}

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 (key.getType() == Boolean.class) {
if (!(value instanceof Boolean)) {
throw invalidValueType(k, value, Boolean.class);
}
b = (Boolean) value;
} else if (key.getType() == String.class) {
if (!(value instanceof String)) {
throw invalidValueType(k, value, String.class);
}
b = FlightSqlSessionUtil.parseStrictBoolean((String) value, name);
} else {
AdbcConnection.super.setOption(key, value);
return;
}
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 (key.getType() == String[].class) {
if (!(value instanceof String[])) {
throw invalidValueType(k, value, String[].class);
}
arr = FlightSqlSessionUtil.validateStringArray((String[]) value);
} else if (key.getType() == String.class) {
if (!(value instanceof String)) {
throw invalidValueType(k, value, String.class);
}
arr = FlightSqlSessionUtil.parseJsonArray((String) value);
} else {
AdbcConnection.super.setOption(key, value);
return;
}
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 (key.getType() == String.class) {
if (!(value instanceof String)) {
throw invalidValueType(k, value, String.class);
}
sv = SessionOptionValueFactory.makeSessionOptionValue((String) value);
} else if (key.getType() == Long.class) {
if (!(value instanceof Long)) {
throw invalidValueType(k, value, Long.class);
}
sv = SessionOptionValueFactory.makeSessionOptionValue((Long) value);
} else if (key.getType() == Double.class) {
if (!(value instanceof Double)) {
throw invalidValueType(k, value, Double.class);
}
sv = SessionOptionValueFactory.makeSessionOptionValue((Double) value);
} else {
AdbcConnection.super.setOption(key, value);
return;
}
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");
}

private static AdbcException invalidValueType(String key, Object value, Class<?> expectedType) {
return AdbcException.invalidArgument(
"[Flight SQL] invalid value type for key "
+ key
+ ": expected "
+ expectedType.getSimpleName()
+ ", got "
+ value.getClass().getSimpleName());
}

@Override
public void close() throws AdbcException {
clientCache.invalidateAll();
try {
AutoCloseables.close(client, allocator);
AutoCloseables.close(
() -> {
try {
// Best-effort: the Go driver also ignores all errors closing the session.
client.closeSession(new CloseSessionRequest());

@Kyperr Kyperr Sep 2, 2026

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Some servers will throw an exception if this does not pass the callOptions. Looks like you went through the effort to implement those options, but they aren't used here.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think it would be worth making the private client gettable() from this class to allow a user to manage sessions while re-using the connection object. Thoughts?

} catch (FlightRuntimeException e) {
// ignore
}
},
clientCache::invalidateAll,
client,
allocator);
} catch (Exception e) {
throw AdbcException.internal("[Flight SQL] Failed to close connection").withCause(e);
}
Expand All @@ -223,6 +395,39 @@
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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,4 +35,10 @@ public interface FlightSqlConnectionProperties {
TypedKey<Boolean> WITH_COOKIE_MIDDLEWARE =
new TypedKey<>("adbc.flight.sql.rpc.with_cookie_middleware", Boolean.class);
String RPC_CALL_HEADER_PREFIX = "adbc.flight.sql.rpc.call_header.";

String SESSION_OPTIONS = "adbc.flight.sql.session.options";
String SESSION_OPTION_PREFIX = "adbc.flight.sql.session.option.";
String SESSION_OPTION_BOOL_PREFIX = "adbc.flight.sql.session.optionbool.";
String SESSION_OPTION_STRING_LIST_PREFIX = "adbc.flight.sql.session.optionstringlist.";
String SESSION_OPTION_ERASE_PREFIX = "adbc.flight.sql.session.optionerase.";
}
Loading
Loading