This is an automated email from the ASF dual-hosted git repository.
lidavidm pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-adbc.git
The following commit(s) were added to refs/heads/main by this push:
new a3a94d4b7 fix(java/driver/flight-sql): include connectionOptions for
preparedStatement.close() (#4513)
a3a94d4b7 is described below
commit a3a94d4b79c145f493801e8517d808a8d8a6ebd4
Author: Daniel_McBride <[email protected]>
AuthorDate: Sat Jul 18 18:24:25 2026 -0700
fix(java/driver/flight-sql): include connectionOptions for
preparedStatement.close() (#4513)
Proposed solution for #4512
Include connection options when calling preparedStatement.close()
---------
Co-authored-by: dmcbride <[email protected]>
---
.../adbc/driver/flightsql/FlightSqlConnection.java | 4 +--
.../adbc/driver/flightsql/FlightSqlStatement.java | 16 +++++----
.../arrow/adbc/driver/flightsql/HeaderTest.java | 39 +++++++++++++++++++++-
.../adbc/driver/flightsql/HeaderValidator.java | 4 +++
4 files changed, 54 insertions(+), 9 deletions(-)
diff --git
a/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlConnection.java
b/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlConnection.java
index 123b9828b..99c9fb735 100644
---
a/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlConnection.java
+++
b/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlConnection.java
@@ -126,7 +126,7 @@ public class FlightSqlConnection implements AdbcConnection {
@Override
public AdbcStatement createStatement() throws AdbcException {
- return new FlightSqlStatement(allocator, client, clientCache, quirks);
+ return new FlightSqlStatement(allocator, client, clientCache, quirks,
callOptions);
}
@Override
@@ -157,7 +157,7 @@ public class FlightSqlConnection implements AdbcConnection {
public AdbcStatement bulkIngest(String targetTableName, BulkIngestMode mode)
throws AdbcException {
return FlightSqlStatement.ingestRoot(
- allocator, client, clientCache, quirks, targetTableName, mode);
+ allocator, client, clientCache, quirks, targetTableName, mode,
callOptions);
}
@Override
diff --git
a/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlStatement.java
b/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlStatement.java
index dd3ba1f28..557353c74 100644
---
a/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlStatement.java
+++
b/java/driver/flight-sql/src/main/java/org/apache/arrow/adbc/driver/flightsql/FlightSqlStatement.java
@@ -29,6 +29,7 @@ import org.apache.arrow.adbc.core.AdbcStatusCode;
import org.apache.arrow.adbc.core.BulkIngestMode;
import org.apache.arrow.adbc.core.PartitionDescriptor;
import org.apache.arrow.adbc.sql.SqlQuirks;
+import org.apache.arrow.flight.CallOption;
import org.apache.arrow.flight.FlightEndpoint;
import org.apache.arrow.flight.FlightInfo;
import org.apache.arrow.flight.FlightRuntimeException;
@@ -36,7 +37,6 @@ import org.apache.arrow.flight.Location;
import org.apache.arrow.flight.impl.Flight;
import org.apache.arrow.flight.sql.FlightSqlClient;
import org.apache.arrow.memory.BufferAllocator;
-import org.apache.arrow.util.AutoCloseables;
import org.apache.arrow.vector.VectorSchemaRoot;
import org.apache.arrow.vector.types.pojo.Field;
import org.apache.arrow.vector.types.pojo.Schema;
@@ -47,6 +47,7 @@ public class FlightSqlStatement implements AdbcStatement {
private final FlightSqlClientWithCallOptions client;
private final LoadingCache<Location, FlightSqlClientWithCallOptions>
clientCache;
private final SqlQuirks quirks;
+ private final CallOption[] connectionOptions;
// State for SQL queries
private @Nullable String sqlQuery;
@@ -59,7 +60,8 @@ public class FlightSqlStatement implements AdbcStatement {
BufferAllocator allocator,
FlightSqlClientWithCallOptions client,
LoadingCache<Location, FlightSqlClientWithCallOptions> clientCache,
- SqlQuirks quirks) {
+ SqlQuirks quirks,
+ CallOption... connectionOptions) {
this.allocator = allocator;
this.client = client;
this.clientCache = clientCache;
@@ -68,6 +70,7 @@ public class FlightSqlStatement implements AdbcStatement {
this.preparedStatement = null;
this.bulkOperation = null;
this.bindRoot = null;
+ this.connectionOptions = connectionOptions;
}
static FlightSqlStatement ingestRoot(
@@ -76,10 +79,11 @@ public class FlightSqlStatement implements AdbcStatement {
LoadingCache<Location, FlightSqlClientWithCallOptions> clientCache,
SqlQuirks quirks,
String targetTableName,
- BulkIngestMode mode) {
+ BulkIngestMode mode,
+ CallOption... connectionOptions) {
Objects.requireNonNull(targetTableName);
final FlightSqlStatement statement =
- new FlightSqlStatement(allocator, client, clientCache, quirks);
+ new FlightSqlStatement(allocator, client, clientCache, quirks,
connectionOptions);
statement.bulkOperation = new BulkState(mode, targetTableName);
return statement;
}
@@ -170,7 +174,7 @@ public class FlightSqlStatement implements AdbcStatement {
statement.setParameters(new NonOwningRoot(bindParams));
client.executePreparedUpdate(statement);
} finally {
- statement.close();
+ statement.close(connectionOptions);
}
} catch (FlightRuntimeException e) {
// XXX: FlightSqlClient.executeUpdate does some extra wrapping that we
need to undo
@@ -313,7 +317,7 @@ public class FlightSqlStatement implements AdbcStatement {
// TODO(https://github.com/apache/arrow/issues/39814): this is annotated
wrongly upstream
if (preparedStatement != null) {
try {
- AutoCloseables.close(preparedStatement);
+ preparedStatement.close(connectionOptions);
} catch (Exception e) {
throw AdbcException.internal("[Flight SQL] Could not close prepared
statement")
.withCause(e);
diff --git
a/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderTest.java
b/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderTest.java
index 261d172e8..b0a25758a 100644
---
a/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderTest.java
+++
b/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderTest.java
@@ -27,6 +27,7 @@ import org.apache.arrow.adbc.core.AdbcDatabase;
import org.apache.arrow.adbc.core.AdbcDriver;
import org.apache.arrow.adbc.core.AdbcException;
import org.apache.arrow.adbc.core.AdbcInfoCode;
+import org.apache.arrow.adbc.core.AdbcStatement;
import org.apache.arrow.adbc.core.AdbcStatusCode;
import org.apache.arrow.adbc.drivermanager.AdbcDriverManager;
import org.apache.arrow.driver.jdbc.utils.MockFlightSqlProducer;
@@ -55,16 +56,18 @@ public class HeaderTest {
private AdbcConnection connection;
private BufferAllocator allocator;
private HeaderValidator.Factory headerValidatorFactory;
+ private MockFlightSqlProducer producer;
@BeforeEach
public void setUp() {
allocator = new RootAllocator(Long.MAX_VALUE);
headerValidatorFactory = new HeaderValidator.Factory();
+ producer = new MockFlightSqlProducer();
builder =
FlightServer.builder()
.middleware(HeaderValidator.KEY, headerValidatorFactory)
.location(Location.forGrpcInsecure("localhost", 0))
- .producer(new MockFlightSqlProducer());
+ .producer(producer);
params = new HashMap<>();
}
@@ -89,6 +92,40 @@ public class HeaderTest {
assertEquals(dummyValue, headers.get(dummyHeaderName));
}
+ /**
+ * The connection's call options (which carry arbitrary headers, auth,
cookies, etc.) must ride on
+ * every RPC the driver issues on that connection's behalf, including the
ClosePreparedStatement
+ * call made when a prepared statement is closed. This regression test
guards against dropping the
+ * connection options on close.
+ */
+ @Test
+ public void testHeaderSentWhenClosingPreparedStatement() throws Exception {
+ final String dummyValue = "dummy";
+ final String dummyHeaderName = "test-header";
+ final String query = "UPDATE the_table SET x = 1";
+ params.put(FlightSqlConnectionProperties.RPC_CALL_HEADER_PREFIX +
dummyHeaderName, dummyValue);
+ producer.addUpdateQuery(query, /* updatedRows= */ 1L);
+ server = builder.build();
+ server.start();
+ connect();
+
+ // Prepare and then close a statement. Closing triggers a
ClosePreparedStatement RPC, which must
+ // still carry the connection header. Before the fix, close() dropped the
call options.
+ try (AdbcStatement statement = connection.createStatement()) {
+ statement.setSqlQuery(query);
+ statement.prepare();
+ }
+
+ // The connection header must appear on every RPC, including the final
ClosePreparedStatement.
+ final int requestCount = headerValidatorFactory.getRequestCount();
+ assertTrue(requestCount > 0);
+ for (int i = 0; i < requestCount; i++) {
+ CallHeaders headers =
headerValidatorFactory.getHeadersReceivedAtRequest(i);
+ assertEquals(
+ dummyValue, headers.get(dummyHeaderName), "connection header missing
on RPC #" + i);
+ }
+ }
+
@Test
public void testCookies() throws Exception {
builder.middleware(CookieMiddleware.KEY, new CookieMiddleware.Factory());
diff --git
a/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderValidator.java
b/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderValidator.java
index f543dd995..c4702eeed 100644
---
a/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderValidator.java
+++
b/java/driver/flight-sql/src/test/java/org/apache/arrow/adbc/driver/flightsql/HeaderValidator.java
@@ -52,6 +52,10 @@ public class HeaderValidator implements
FlightServerMiddleware {
return cloneHeaders(headersReceived.get(request));
}
+ public int getRequestCount() {
+ return headersReceived.size();
+ }
+
private static CallHeaders cloneHeaders(CallHeaders headers) {
FlightCallHeaders cloneHeaders = new FlightCallHeaders();
for (String key : headers.keys()) {