diff --git a/src/main/java/com/databricks/jdbc/api/impl/BatchParameterSet.java b/src/main/java/com/databricks/jdbc/api/impl/BatchParameterSet.java
new file mode 100644
index 000000000..6edc670db
--- /dev/null
+++ b/src/main/java/com/databricks/jdbc/api/impl/BatchParameterSet.java
@@ -0,0 +1,88 @@
+package com.databricks.jdbc.api.impl;
+
+import java.sql.Date;
+import java.sql.Time;
+import java.sql.Timestamp;
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+/**
+ * Immutable, position-ordered snapshot of one prepared-statement parameter set.
+ *
+ *
This model preserves JDBC's one-based parameter indexes. Transport adapters are responsible
+ * for converting them to protocol-specific wire ordinals. It does not validate parameter
+ * completeness, index continuity, or consistency with other parameter sets; those validations
+ * remain the backend's responsibility.
+ */
+public final class BatchParameterSet {
+
+ private final List parameters;
+ private final Map parameterBindings;
+
+ private BatchParameterSet(List parameters) {
+ this.parameters = List.copyOf(parameters);
+ Map bindings = new LinkedHashMap<>();
+ this.parameters.forEach(parameter -> bindings.put(parameter.cardinal(), parameter));
+ this.parameterBindings = Collections.unmodifiableMap(bindings);
+ }
+
+ public static BatchParameterSet from(Map parameterBindings) {
+ Objects.requireNonNull(parameterBindings, "parameterBindings");
+ List orderedParameters =
+ parameterBindings.entrySet().stream()
+ .sorted(Comparator.comparingInt(Map.Entry::getKey))
+ .map(BatchParameterSet::snapshotParameter)
+ .collect(Collectors.toList());
+ return new BatchParameterSet(orderedParameters);
+ }
+
+ public List getParameters() {
+ return parameters;
+ }
+
+ public Map getParameterBindings() {
+ return parameterBindings;
+ }
+
+ public int size() {
+ return parameters.size();
+ }
+
+ public boolean isEmpty() {
+ return parameters.isEmpty();
+ }
+
+ private static ImmutableSqlParameter snapshotParameter(
+ Map.Entry entry) {
+ ImmutableSqlParameter parameter = entry.getValue();
+ return ImmutableSqlParameter.builder()
+ .cardinal(entry.getKey())
+ .type(parameter.type())
+ .value(snapshotValue(parameter.value()))
+ .build();
+ }
+
+ private static Object snapshotValue(Object value) {
+ if (value instanceof Timestamp) {
+ Timestamp timestamp = (Timestamp) value;
+ Timestamp copy = new Timestamp(timestamp.getTime());
+ copy.setNanos(timestamp.getNanos());
+ return copy;
+ }
+ if (value instanceof Date) {
+ return new Date(((Date) value).getTime());
+ }
+ if (value instanceof Time) {
+ return new Time(((Time) value).getTime());
+ }
+ if (value instanceof byte[]) {
+ return ((byte[]) value).clone();
+ }
+ return value;
+ }
+}
diff --git a/src/main/java/com/databricks/jdbc/api/impl/DatabricksConnectionContext.java b/src/main/java/com/databricks/jdbc/api/impl/DatabricksConnectionContext.java
index dfa4b70f9..62233e70c 100644
--- a/src/main/java/com/databricks/jdbc/api/impl/DatabricksConnectionContext.java
+++ b/src/main/java/com/databricks/jdbc/api/impl/DatabricksConnectionContext.java
@@ -1480,6 +1480,11 @@ public boolean isBatchedInsertsEnabled() {
return getParameter(DatabricksJdbcUrlParams.ENABLE_BATCHED_INSERTS).equals("1");
}
+ @Override
+ public boolean isNativeBatchingEnabled() {
+ return getParameter(DatabricksJdbcUrlParams.ENABLE_NATIVE_BATCHING).equals("1");
+ }
+
@Override
public List getNonRowcountQueryPrefixes() {
String prefixesStr = getParameter(DatabricksJdbcUrlParams.NON_ROWCOUNT_QUERY_PREFIXES);
diff --git a/src/main/java/com/databricks/jdbc/api/impl/DatabricksPreparedStatement.java b/src/main/java/com/databricks/jdbc/api/impl/DatabricksPreparedStatement.java
index 32e856aac..27aea28d6 100644
--- a/src/main/java/com/databricks/jdbc/api/impl/DatabricksPreparedStatement.java
+++ b/src/main/java/com/databricks/jdbc/api/impl/DatabricksPreparedStatement.java
@@ -33,7 +33,7 @@ public class DatabricksPreparedStatement extends DatabricksStatement implements
JdbcLoggerFactory.getLogger(DatabricksPreparedStatement.class);
private final String sql;
private DatabricksParameterMetaData databricksParameterMetaData;
- private List databricksBatchParameterMetaData;
+ private List batchParameterSets;
private final boolean interpolateParameters;
private final int CHUNK_SIZE = 8192;
@@ -43,7 +43,7 @@ public DatabricksPreparedStatement(DatabricksConnection connection, String sql)
this.sql = sql;
this.interpolateParameters = connection.getConnectionContext().supportManyParameters();
this.databricksParameterMetaData = new DatabricksParameterMetaData(sql);
- this.databricksBatchParameterMetaData = new ArrayList<>();
+ this.batchParameterSets = new ArrayList<>();
// Cache whether this statement should return a ResultSet (based on SQL and config)
this.shouldReturnResultSet = shouldReturnResultSetWithConfig(sql);
}
@@ -58,7 +58,7 @@ public DatabricksPreparedStatement(DatabricksConnection connection, String sql)
this.sql = sql;
this.interpolateParameters = interpolateParameters;
this.databricksParameterMetaData = databricksParameterMetaData;
- this.databricksBatchParameterMetaData = new ArrayList<>();
+ this.batchParameterSets = new ArrayList<>();
// Cache whether this statement should return a ResultSet (based on SQL and config)
this.shouldReturnResultSet = shouldReturnResultSetWithConfig(sql);
}
@@ -94,7 +94,7 @@ public int executeUpdate() throws SQLException {
}
@Override
- public int[] executeBatch() throws DatabricksBatchUpdateException {
+ public int[] executeBatch() throws SQLException {
LOGGER.debug("public int executeBatch()");
long[] largeUpdateCount = executeLargeBatch();
int[] updateCount = new int[largeUpdateCount.length];
@@ -107,10 +107,10 @@ public int[] executeBatch() throws DatabricksBatchUpdateException {
}
@Override
- public long[] executeLargeBatch() throws DatabricksBatchUpdateException {
+ public long[] executeLargeBatch() throws SQLException {
LOGGER.debug("public long executeLargeBatch()");
- if (databricksBatchParameterMetaData.isEmpty()) {
+ if (batchParameterSets.isEmpty()) {
return new long[0];
}
@@ -121,18 +121,41 @@ public long[] executeLargeBatch() throws DatabricksBatchUpdateException {
connection,
interpolateParameters,
(sqlToExecute, params, statementType, closeStatement) ->
- executeInternal(sqlToExecute, params, statementType, closeStatement));
-
- long[] updateCounts = batchExecutor.executeBatch(databricksBatchParameterMetaData);
+ executeInternal(sqlToExecute, params, statementType, closeStatement),
+ new PreparedStatementBatchExecutor.NativeBatchExecutor() {
+ @Override
+ public boolean isSupported() {
+ return supportsNativeParameterBatching();
+ }
+
+ @Override
+ public long[] execute(String sql, List parameterSets)
+ throws SQLException {
+ return executeNativeBatchInternal(sql, parameterSets);
+ }
+ });
+
+ long[] updateCounts;
+ try {
+ updateCounts = batchExecutor.executeBatch(batchParameterSets);
+ } catch (NativeBatchResultException e) {
+ // The backend already completed the batch. Clear it before propagating the count-read error
+ // so a caller retry cannot insert the same rows again.
+ clearBatchAfterExecution();
+ throw e;
+ }
// Clear the batch after successful execution per JDBC spec
+ clearBatchAfterExecution();
+ return updateCounts;
+ }
+
+ private void clearBatchAfterExecution() {
try {
clearBatch();
} catch (SQLException e) {
- LOGGER.error("Failed to clear batch after successful execution", e);
+ LOGGER.error("Failed to clear batch after execution", e);
}
-
- return updateCounts;
}
@Override
@@ -371,7 +394,8 @@ public boolean execute() throws SQLException {
@Override
public void addBatch() {
LOGGER.debug("public void addBatch()");
- this.databricksBatchParameterMetaData.add(databricksParameterMetaData);
+ this.batchParameterSets.add(
+ BatchParameterSet.from(databricksParameterMetaData.getParameterBindings()));
this.databricksParameterMetaData = new DatabricksParameterMetaData(sql);
}
@@ -380,7 +404,7 @@ public void clearBatch() throws DatabricksSQLException {
LOGGER.debug("public void clearBatch()");
checkIfClosed();
this.databricksParameterMetaData = new DatabricksParameterMetaData(sql);
- this.databricksBatchParameterMetaData = new ArrayList<>();
+ this.batchParameterSets = new ArrayList<>();
}
@Override
@@ -755,7 +779,7 @@ private void checkLength(long targetLength, long sourceLength) throws SQLExcepti
}
private void checkIfBatchOperation() throws DatabricksSQLException {
- if (!this.databricksBatchParameterMetaData.isEmpty()) {
+ if (!this.batchParameterSets.isEmpty()) {
String errorMessage =
"Batch must either be executed with executeBatch() or cleared with clearBatch()";
LOGGER.error(errorMessage);
diff --git a/src/main/java/com/databricks/jdbc/api/impl/DatabricksResultSet.java b/src/main/java/com/databricks/jdbc/api/impl/DatabricksResultSet.java
index cde481ccf..e0fbb4ccf 100644
--- a/src/main/java/com/databricks/jdbc/api/impl/DatabricksResultSet.java
+++ b/src/main/java/com/databricks/jdbc/api/impl/DatabricksResultSet.java
@@ -62,6 +62,7 @@ enum ResultSetType {
private static final JdbcLogger LOGGER = JdbcLoggerFactory.getLogger(DatabricksResultSet.class);
protected static final String AFFECTED_ROWS_COUNT = "num_affected_rows";
+ private static final String REPEAT_COUNT = "repeat";
private final ExecutionStatus executionStatus;
private final StatementId statementId;
private final IExecutionResult executionResult;
@@ -2310,6 +2311,44 @@ public long getUpdateCount() throws SQLException {
return updateCount;
}
+ long[] getBatchUpdateCounts(int expectedCount) throws SQLException {
+ checkIfClosed();
+ if (resultSetMetaData.getColumnNameIndex(AFFECTED_ROWS_COUNT) < 1) {
+ throw new DatabricksSQLException(
+ "Native batch result is missing column " + AFFECTED_ROWS_COUNT,
+ DatabricksDriverErrorCode.RESULT_SET_ERROR);
+ }
+
+ long[] counts = new long[expectedCount];
+ int index = 0;
+ boolean hasRepeatCount = resultSetMetaData.getColumnNameIndex(REPEAT_COUNT) > 0;
+ countingUpdateRows = true;
+ try {
+ while (next()) {
+ long repeatCount = hasRepeatCount ? getLong(REPEAT_COUNT) : 1;
+ if (repeatCount < 1 || repeatCount > expectedCount - index) {
+ throw new DatabricksSQLException(
+ "Native batch returned an invalid repeat count: " + repeatCount,
+ DatabricksDriverErrorCode.RESULT_SET_ERROR);
+ }
+ long affectedRows = getLong(AFFECTED_ROWS_COUNT);
+ for (long repeated = 0; repeated < repeatCount; repeated++) {
+ counts[index++] = affectedRows;
+ }
+ }
+ } finally {
+ countingUpdateRows = false;
+ }
+
+ if (index != expectedCount) {
+ throw new DatabricksSQLException(
+ String.format(
+ "Native batch returned %d update counts for %d parameter sets", index, expectedCount),
+ DatabricksDriverErrorCode.RESULT_SET_ERROR);
+ }
+ return counts;
+ }
+
@Override
public boolean hasUpdateCount() throws SQLException {
checkIfClosed();
diff --git a/src/main/java/com/databricks/jdbc/api/impl/DatabricksStatement.java b/src/main/java/com/databricks/jdbc/api/impl/DatabricksStatement.java
index d1dd0d30e..f9493c50b 100644
--- a/src/main/java/com/databricks/jdbc/api/impl/DatabricksStatement.java
+++ b/src/main/java/com/databricks/jdbc/api/impl/DatabricksStatement.java
@@ -866,6 +866,15 @@ DatabricksResultSet executeInternal(
LOGGER.debug(stackTraceMessage);
CompletableFuture futureResultSet =
getFutureResult(sql, params, statementType);
+ return waitForExecutionResult(sql, stackTraceMessage, futureResultSet, closeStatement);
+ }
+
+ private DatabricksResultSet waitForExecutionResult(
+ String sql,
+ String stackTraceMessage,
+ CompletableFuture futureResultSet,
+ boolean closeStatement)
+ throws SQLException {
try {
resultSet =
timeoutInSeconds == 0
@@ -938,6 +947,38 @@ DatabricksResultSet executeInternal(
return result;
}
+ boolean supportsNativeParameterBatching() {
+ try {
+ IDatabricksClient client = connection.getSession().getDatabricksClient();
+ return client.supportsNativeParameterBatching(connection.getSession().getComputeResource());
+ } catch (DatabricksSQLException e) {
+ LOGGER.warn("Unable to determine native batch capability, using legacy execution", e);
+ return false;
+ }
+ }
+
+ long[] executeNativeBatchInternal(String sql, List parameterSets)
+ throws SQLException {
+ resetForNewExecution();
+ DatabricksThreadContextHolder.setStatementType(StatementType.UPDATE);
+ String stackTraceMessage =
+ format(
+ "DatabricksResultSet executeNativeBatchInternal(String sql = %s, parameterSetCount = %s)",
+ sql, parameterSets.size());
+ LOGGER.debug(stackTraceMessage);
+ DatabricksResultSet result =
+ waitForExecutionResult(
+ sql,
+ stackTraceMessage,
+ getFutureBatchResult(sql, parameterSets, StatementType.UPDATE),
+ true);
+ try {
+ return result.getBatchUpdateCounts(parameterSets.size());
+ } catch (SQLException e) {
+ throw new NativeBatchResultException(e);
+ }
+ }
+
CompletableFuture getFutureResult(
String sql, Map params, StatementType statementType) {
return CompletableFuture.supplyAsync(
@@ -954,6 +995,21 @@ CompletableFuture getFutureResult(
executor);
}
+ private CompletableFuture getFutureBatchResult(
+ String sql, List parameterSets, StatementType statementType) {
+ return CompletableFuture.supplyAsync(
+ () -> {
+ try {
+ String sqlString = escapeProcessing ? StringUtil.convertJdbcEscapeSequences(sql) : sql;
+ sqlString = StringUtil.removeRedundantEscapeClause(sqlString);
+ return getBatchResultFromClient(sqlString, parameterSets, statementType);
+ } catch (SQLException e) {
+ throw new RuntimeException(e);
+ }
+ },
+ executor);
+ }
+
DatabricksResultSet getResultFromClient(
String sql, Map params, StatementType statementType)
throws SQLException {
@@ -968,6 +1024,19 @@ DatabricksResultSet getResultFromClient(
null /* metadataOperationType */);
}
+ private DatabricksResultSet getBatchResultFromClient(
+ String sql, List parameterSets, StatementType statementType)
+ throws SQLException {
+ IDatabricksClient client = connection.getSession().getDatabricksClient();
+ return client.executeStatementBatch(
+ sql,
+ connection.getSession().getComputeResource(),
+ parameterSets,
+ statementType,
+ connection.getSession(),
+ this);
+ }
+
void checkIfClosed() throws DatabricksSQLException {
if (isClosed) {
throw new DatabricksSQLException(
diff --git a/src/main/java/com/databricks/jdbc/api/impl/LegacyPreparedStatementBatchExecutor.java b/src/main/java/com/databricks/jdbc/api/impl/LegacyPreparedStatementBatchExecutor.java
new file mode 100644
index 000000000..05663b1c9
--- /dev/null
+++ b/src/main/java/com/databricks/jdbc/api/impl/LegacyPreparedStatementBatchExecutor.java
@@ -0,0 +1,207 @@
+package com.databricks.jdbc.api.impl;
+
+import com.databricks.jdbc.common.DatabricksJdbcConstants;
+import com.databricks.jdbc.common.StatementType;
+import com.databricks.jdbc.common.util.InsertStatementParser;
+import com.databricks.jdbc.exception.DatabricksBatchUpdateException;
+import com.databricks.jdbc.exception.DatabricksSQLException;
+import com.databricks.jdbc.log.JdbcLogger;
+import com.databricks.jdbc.log.JdbcLoggerFactory;
+import com.databricks.jdbc.model.telemetry.enums.DatabricksDriverErrorCode;
+import java.sql.Statement;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * Executes prepared-statement batches using the legacy client-side strategies.
+ *
+ * This class intentionally preserves the existing behavior: eligible INSERT statements may be
+ * rewritten into chunked multi-row INSERTs, while all other statements execute one parameter set at
+ * a time.
+ */
+class LegacyPreparedStatementBatchExecutor {
+
+ private static final JdbcLogger LOGGER =
+ JdbcLoggerFactory.getLogger(LegacyPreparedStatementBatchExecutor.class);
+
+ private final String sql;
+ private final DatabricksConnection connection;
+ private final boolean interpolateParameters;
+ private final PreparedStatementBatchExecutor.StatementExecutor statementExecutor;
+
+ LegacyPreparedStatementBatchExecutor(
+ String sql,
+ DatabricksConnection connection,
+ boolean interpolateParameters,
+ PreparedStatementBatchExecutor.StatementExecutor statementExecutor) {
+ this.sql = sql;
+ this.connection = connection;
+ this.interpolateParameters = interpolateParameters;
+ this.statementExecutor = statementExecutor;
+ }
+
+ long[] executeBatch(List batchParameterSets)
+ throws DatabricksBatchUpdateException {
+ if (batchParameterSets.isEmpty()) {
+ return new long[0];
+ }
+
+ // Try to optimize INSERT statements with multi-row batching
+ if (canUseBatchedInsert()) {
+ return executeBatchedInsert(batchParameterSets);
+ } else {
+ // Fall back to individual execution for non-INSERT or incompatible statements
+ return executeIndividualStatements(batchParameterSets);
+ }
+ }
+
+ long[] executeIndividually(List batchParameterSets)
+ throws DatabricksBatchUpdateException {
+ return executeIndividualStatements(batchParameterSets);
+ }
+
+ private boolean canUseBatchedInsert() {
+ // Check if batched inserts are enabled via connection property
+ if (!connection.getConnectionContext().isBatchedInsertsEnabled()) {
+ return false;
+ }
+
+ // Use strict exception-based parsing for better error handling
+ try {
+ InsertStatementParser.parseInsertStrict(sql);
+ return true;
+ } catch (Exception e) {
+ // Not a valid INSERT statement suitable for batching
+ LOGGER.warn(
+ "EnableBatchedInserts is enabled but the INSERT statement could not be parsed for"
+ + " batching, falling back to individual execution: {}",
+ e.getMessage());
+ return false;
+ }
+ }
+
+ private long[] executeBatchedInsert(List batchParameterSets)
+ throws DatabricksBatchUpdateException {
+ LOGGER.debug("Executing batched INSERT with {} rows", batchParameterSets.size());
+
+ try {
+ InsertStatementParser.InsertInfo insertInfo = InsertStatementParser.parseInsertStrict(sql);
+
+ // Calculate how many rows we can fit in one chunk
+ int parametersPerRow = insertInfo.getColumnCount();
+ int maxRowsPerChunk;
+
+ if (interpolateParameters) {
+ // When parameter interpolation is enabled (supportManyParameters=1), there is no
+ // parameter limit since values are interpolated directly into the SQL string.
+ // Try to execute all rows in a single batch, only limited by configured BatchInsertSize
+ // which users can set based on their data to avoid exceeding the 16MB statement limit.
+ int configuredBatchSize = connection.getConnectionContext().getBatchInsertSize();
+ if (configuredBatchSize < 1) {
+ throw new DatabricksSQLException(
+ "BatchInsertSize must be at least 1, got: " + configuredBatchSize,
+ DatabricksDriverErrorCode.INVALID_STATE);
+ }
+ maxRowsPerChunk = Math.min(configuredBatchSize, batchParameterSets.size());
+ } else {
+ // When using parameterized queries, respect the 256 parameter limit from Databricks
+ // backend
+ int maxRowsByParameterLimit =
+ DatabricksJdbcConstants.MAX_QUERY_PARAMETERS / parametersPerRow;
+
+ // Ensure we have at least 1 row per chunk
+ if (maxRowsByParameterLimit < 1) {
+ maxRowsPerChunk = 1;
+ } else {
+ maxRowsPerChunk = maxRowsByParameterLimit;
+ }
+ }
+
+ long[] allUpdateCounts = new long[batchParameterSets.size()];
+
+ // Process batches in chunks
+ for (int startIndex = 0;
+ startIndex < batchParameterSets.size();
+ startIndex += maxRowsPerChunk) {
+ int endIndex = Math.min(startIndex + maxRowsPerChunk, batchParameterSets.size());
+ int chunkSize = endIndex - startIndex;
+
+ // Build multi-row INSERT for this chunk
+ String multiRowSql = InsertStatementParser.generateMultiRowInsert(insertInfo, chunkSize);
+ Map chunkParams = new HashMap<>();
+ int paramIndex = 1;
+
+ for (int i = startIndex; i < endIndex; i++) {
+ BatchParameterSet batchParams = batchParameterSets.get(i);
+ Map rowParams = batchParams.getParameterBindings();
+ for (int j = 1; j <= rowParams.size(); j++) {
+ if (rowParams.containsKey(j)) {
+ chunkParams.put(paramIndex++, rowParams.get(j));
+ }
+ }
+ }
+
+ // Execute this chunk
+ String sqlToExecute =
+ interpolateParameters
+ ? com.databricks.jdbc.common.util.SQLInterpolator.interpolateSQL(
+ multiRowSql, chunkParams)
+ : multiRowSql;
+ Map paramsToSend =
+ interpolateParameters ? new HashMap<>() : chunkParams;
+ statementExecutor.execute(sqlToExecute, paramsToSend, StatementType.UPDATE, false);
+
+ // Set update counts for this chunk (each row typically affects 1 row)
+ for (int i = startIndex; i < endIndex; i++) {
+ allUpdateCounts[i] = 1;
+ }
+ }
+
+ return allUpdateCounts;
+
+ } catch (DatabricksBatchUpdateException e) {
+ // Re-throw batch update exceptions (these already have proper update counts)
+ throw e;
+ } catch (Exception e) {
+ // Unexpected exception - mark all as failed
+ LOGGER.error("Unexpected error executing batched INSERT: {}", e.getMessage(), e);
+ long[] failedCounts = new long[batchParameterSets.size()];
+ for (int i = 0; i < failedCounts.length; i++) {
+ failedCounts[i] = Statement.EXECUTE_FAILED;
+ }
+ throw new DatabricksBatchUpdateException(
+ e.getMessage(), DatabricksDriverErrorCode.BATCH_EXECUTE_EXCEPTION, failedCounts);
+ }
+ }
+
+ private long[] executeIndividualStatements(List batchParameterSets)
+ throws DatabricksBatchUpdateException {
+ LOGGER.debug("Executing batch individually with {} statements", batchParameterSets.size());
+ long[] largeUpdateCount = new long[batchParameterSets.size()];
+
+ for (int sqlQueryIndex = 0; sqlQueryIndex < batchParameterSets.size(); sqlQueryIndex++) {
+ BatchParameterSet batchParameterSet = batchParameterSets.get(sqlQueryIndex);
+ try {
+ DatabricksResultSet resultSet =
+ statementExecutor.execute(
+ sql, batchParameterSet.getParameterBindings(), StatementType.UPDATE, false);
+ largeUpdateCount[sqlQueryIndex] = resultSet.getUpdateCount();
+ } catch (Exception e) {
+ LOGGER.error(
+ "Error executing batch update for index {}: {}", sqlQueryIndex, e.getMessage(), e);
+ // Set the current failed statement's count
+ largeUpdateCount[sqlQueryIndex] = Statement.EXECUTE_FAILED;
+ // Set all remaining statements as failed
+ for (int i = sqlQueryIndex + 1; i < largeUpdateCount.length; i++) {
+ largeUpdateCount[i] = Statement.EXECUTE_FAILED;
+ }
+ // WARNING: Due to lack of transaction support, any successfully executed statements
+ // before this failure have already been committed and cannot be rolled back
+ throw new DatabricksBatchUpdateException(
+ e.getMessage(), DatabricksDriverErrorCode.BATCH_EXECUTE_EXCEPTION, largeUpdateCount);
+ }
+ }
+ return largeUpdateCount;
+ }
+}
diff --git a/src/main/java/com/databricks/jdbc/api/impl/NativeBatchResultException.java b/src/main/java/com/databricks/jdbc/api/impl/NativeBatchResultException.java
new file mode 100644
index 000000000..e121b0bb3
--- /dev/null
+++ b/src/main/java/com/databricks/jdbc/api/impl/NativeBatchResultException.java
@@ -0,0 +1,22 @@
+package com.databricks.jdbc.api.impl;
+
+import com.databricks.jdbc.exception.DatabricksSQLException;
+import com.databricks.jdbc.model.telemetry.enums.DatabricksDriverErrorCode;
+import java.sql.SQLException;
+
+/**
+ * Indicates that a native batch succeeded but its JDBC update counts could not be read.
+ *
+ * This is intentionally not a {@code BatchUpdateException}: backend execution did not fail.
+ */
+class NativeBatchResultException extends DatabricksSQLException {
+
+ NativeBatchResultException(SQLException cause) {
+ super(
+ "Native batch execution succeeded, but JDBC update counts could not be read. "
+ + "Inserted rows may already be committed. Cause: "
+ + cause.getMessage(),
+ cause,
+ DatabricksDriverErrorCode.RESULT_SET_ERROR);
+ }
+}
diff --git a/src/main/java/com/databricks/jdbc/api/impl/PreparedStatementBatchExecutor.java b/src/main/java/com/databricks/jdbc/api/impl/PreparedStatementBatchExecutor.java
index 7387cf67f..3247eecc2 100644
--- a/src/main/java/com/databricks/jdbc/api/impl/PreparedStatementBatchExecutor.java
+++ b/src/main/java/com/databricks/jdbc/api/impl/PreparedStatementBatchExecutor.java
@@ -1,28 +1,33 @@
package com.databricks.jdbc.api.impl;
-import com.databricks.jdbc.common.DatabricksJdbcConstants;
import com.databricks.jdbc.common.StatementType;
import com.databricks.jdbc.common.util.InsertStatementParser;
import com.databricks.jdbc.exception.DatabricksBatchUpdateException;
-import com.databricks.jdbc.exception.DatabricksSQLException;
-import com.databricks.jdbc.log.JdbcLogger;
-import com.databricks.jdbc.log.JdbcLoggerFactory;
-import com.databricks.jdbc.model.telemetry.enums.DatabricksDriverErrorCode;
import java.sql.SQLException;
import java.sql.Statement;
-import java.util.HashMap;
+import java.util.Arrays;
import java.util.List;
import java.util.Map;
class PreparedStatementBatchExecutor {
- private static final JdbcLogger LOGGER =
- JdbcLoggerFactory.getLogger(PreparedStatementBatchExecutor.class);
+ private static final NativeBatchExecutor UNSUPPORTED_NATIVE_EXECUTOR =
+ new NativeBatchExecutor() {
+ @Override
+ public boolean isSupported() {
+ return false;
+ }
+
+ @Override
+ public long[] execute(String sql, List parameterSets) {
+ throw new IllegalStateException("Native batch execution is not supported");
+ }
+ };
private final String sql;
private final DatabricksConnection connection;
- private final boolean interpolateParameters;
- private final StatementExecutor statementExecutor;
+ private final LegacyPreparedStatementBatchExecutor legacyExecutor;
+ private final NativeBatchExecutor nativeExecutor;
@FunctionalInterface
interface StatementExecutor {
@@ -34,178 +39,63 @@ DatabricksResultSet execute(
throws SQLException;
}
+ interface NativeBatchExecutor {
+ boolean isSupported();
+
+ long[] execute(String sql, List parameterSets) throws SQLException;
+ }
+
PreparedStatementBatchExecutor(
String sql,
DatabricksConnection connection,
boolean interpolateParameters,
StatementExecutor statementExecutor) {
+ this(sql, connection, interpolateParameters, statementExecutor, UNSUPPORTED_NATIVE_EXECUTOR);
+ }
+
+ PreparedStatementBatchExecutor(
+ String sql,
+ DatabricksConnection connection,
+ boolean interpolateParameters,
+ StatementExecutor statementExecutor,
+ NativeBatchExecutor nativeExecutor) {
this.sql = sql;
this.connection = connection;
- this.interpolateParameters = interpolateParameters;
- this.statementExecutor = statementExecutor;
+ this.legacyExecutor =
+ new LegacyPreparedStatementBatchExecutor(
+ sql, connection, interpolateParameters, statementExecutor);
+ this.nativeExecutor = nativeExecutor;
}
- long[] executeBatch(List batchParameterMetaData)
- throws DatabricksBatchUpdateException {
- if (batchParameterMetaData.isEmpty()) {
+ long[] executeBatch(List batchParameterSets) throws SQLException {
+ if (batchParameterSets.isEmpty()) {
return new long[0];
}
-
- // Try to optimize INSERT statements with multi-row batching
- if (canUseBatchedInsert()) {
- return executeBatchedInsert(batchParameterMetaData);
- } else {
- // Fall back to individual execution for non-INSERT or incompatible statements
- return executeIndividualStatements(batchParameterMetaData);
- }
- }
-
- private boolean canUseBatchedInsert() {
- // Check if batched inserts are enabled via connection property
- if (!connection.getConnectionContext().isBatchedInsertsEnabled()) {
- return false;
+ if (!InsertStatementParser.isParametrizedInsert(sql)) {
+ return legacyExecutor.executeIndividually(batchParameterSets);
}
-
- // Use strict exception-based parsing for better error handling
- try {
- InsertStatementParser.parseInsertStrict(sql);
- return true;
- } catch (Exception e) {
- // Not a valid INSERT statement suitable for batching
- LOGGER.warn(
- "EnableBatchedInserts is enabled but the INSERT statement could not be parsed for"
- + " batching, falling back to individual execution: {}",
- e.getMessage());
- return false;
+ if (!connection.getConnectionContext().isNativeBatchingEnabled()
+ || !nativeExecutor.isSupported()) {
+ return legacyExecutor.executeBatch(batchParameterSets);
}
- }
-
- private long[] executeBatchedInsert(List batchParameterMetaData)
- throws DatabricksBatchUpdateException {
- LOGGER.debug("Executing batched INSERT with {} rows", batchParameterMetaData.size());
-
try {
- InsertStatementParser.InsertInfo insertInfo = InsertStatementParser.parseInsertStrict(sql);
-
- // Calculate how many rows we can fit in one chunk
- int parametersPerRow = insertInfo.getColumnCount();
- int maxRowsPerChunk;
-
- if (interpolateParameters) {
- // When parameter interpolation is enabled (supportManyParameters=1), there is no
- // parameter limit since values are interpolated directly into the SQL string.
- // Try to execute all rows in a single batch, only limited by configured BatchInsertSize
- // which users can set based on their data to avoid exceeding the 16MB statement limit.
- int configuredBatchSize = connection.getConnectionContext().getBatchInsertSize();
- if (configuredBatchSize < 1) {
- throw new DatabricksSQLException(
- "BatchInsertSize must be at least 1, got: " + configuredBatchSize,
- DatabricksDriverErrorCode.INVALID_STATE);
- }
- maxRowsPerChunk = Math.min(configuredBatchSize, batchParameterMetaData.size());
- } else {
- // When using parameterized queries, respect the 256 parameter limit from Databricks
- // backend
- int maxRowsByParameterLimit =
- DatabricksJdbcConstants.MAX_QUERY_PARAMETERS / parametersPerRow;
-
- // Ensure we have at least 1 row per chunk
- if (maxRowsByParameterLimit < 1) {
- maxRowsPerChunk = 1;
- } else {
- maxRowsPerChunk = maxRowsByParameterLimit;
- }
- }
-
- long[] allUpdateCounts = new long[batchParameterMetaData.size()];
-
- // Process batches in chunks
- for (int startIndex = 0;
- startIndex < batchParameterMetaData.size();
- startIndex += maxRowsPerChunk) {
- int endIndex = Math.min(startIndex + maxRowsPerChunk, batchParameterMetaData.size());
- int chunkSize = endIndex - startIndex;
-
- // Build multi-row INSERT for this chunk
- String multiRowSql = InsertStatementParser.generateMultiRowInsert(insertInfo, chunkSize);
- Map chunkParams = new HashMap<>();
- int paramIndex = 1;
-
- for (int i = startIndex; i < endIndex; i++) {
- DatabricksParameterMetaData batchParams = batchParameterMetaData.get(i);
- Map rowParams = batchParams.getParameterBindings();
- for (int j = 1; j <= rowParams.size(); j++) {
- if (rowParams.containsKey(j)) {
- chunkParams.put(paramIndex++, rowParams.get(j));
- }
- }
- }
-
- // Execute this chunk
- String sqlToExecute =
- interpolateParameters
- ? com.databricks.jdbc.common.util.SQLInterpolator.interpolateSQL(
- multiRowSql, chunkParams)
- : multiRowSql;
- Map paramsToSend =
- interpolateParameters ? new HashMap<>() : chunkParams;
- statementExecutor.execute(sqlToExecute, paramsToSend, StatementType.UPDATE, false);
-
- // Set update counts for this chunk (each row typically affects 1 row)
- for (int i = startIndex; i < endIndex; i++) {
- allUpdateCounts[i] = 1;
- }
- }
-
- return allUpdateCounts;
-
- } catch (DatabricksBatchUpdateException e) {
- // Re-throw batch update exceptions (these already have proper update counts)
+ return nativeExecutor.execute(sql, batchParameterSets);
+ } catch (NativeBatchResultException e) {
throw e;
- } catch (Exception e) {
- // Unexpected exception - mark all as failed
- LOGGER.error("Unexpected error executing batched INSERT: {}", e.getMessage(), e);
- long[] failedCounts = new long[batchParameterMetaData.size()];
- for (int i = 0; i < failedCounts.length; i++) {
- failedCounts[i] = Statement.EXECUTE_FAILED;
+ } catch (SQLException e) {
+ if (isUnsupportedNativeBatching(e)) {
+ return legacyExecutor.executeBatch(batchParameterSets);
}
+ long[] failedCounts = new long[batchParameterSets.size()];
+ Arrays.fill(failedCounts, Statement.EXECUTE_FAILED);
throw new DatabricksBatchUpdateException(
- e.getMessage(), DatabricksDriverErrorCode.BATCH_EXECUTE_EXCEPTION, failedCounts);
+ e.getMessage(), e.getSQLState(), e.getErrorCode(), failedCounts, e);
}
}
- private long[] executeIndividualStatements(
- List batchParameterMetaData)
- throws DatabricksBatchUpdateException {
- LOGGER.debug("Executing batch individually with {} statements", batchParameterMetaData.size());
- long[] largeUpdateCount = new long[batchParameterMetaData.size()];
-
- for (int sqlQueryIndex = 0; sqlQueryIndex < batchParameterMetaData.size(); sqlQueryIndex++) {
- DatabricksParameterMetaData databricksParameterMetaData =
- batchParameterMetaData.get(sqlQueryIndex);
- try {
- DatabricksResultSet resultSet =
- statementExecutor.execute(
- sql,
- databricksParameterMetaData.getParameterBindings(),
- StatementType.UPDATE,
- false);
- largeUpdateCount[sqlQueryIndex] = resultSet.getUpdateCount();
- } catch (Exception e) {
- LOGGER.error(
- "Error executing batch update for index {}: {}", sqlQueryIndex, e.getMessage(), e);
- // Set the current failed statement's count
- largeUpdateCount[sqlQueryIndex] = Statement.EXECUTE_FAILED;
- // Set all remaining statements as failed
- for (int i = sqlQueryIndex + 1; i < largeUpdateCount.length; i++) {
- largeUpdateCount[i] = Statement.EXECUTE_FAILED;
- }
- // WARNING: Due to lack of transaction support, any successfully executed statements
- // before this failure have already been committed and cannot be rolled back
- throw new DatabricksBatchUpdateException(
- e.getMessage(), DatabricksDriverErrorCode.BATCH_EXECUTE_EXCEPTION, largeUpdateCount);
- }
- }
- return largeUpdateCount;
+ private boolean isUnsupportedNativeBatching(SQLException exception) {
+ return "42P02".equals(exception.getSQLState())
+ && exception.getMessage() != null
+ && exception.getMessage().contains("[UNBOUND_SQL_PARAMETER]");
}
}
diff --git a/src/main/java/com/databricks/jdbc/api/internal/IDatabricksConnectionContext.java b/src/main/java/com/databricks/jdbc/api/internal/IDatabricksConnectionContext.java
index fb0d745a7..c3851296c 100644
--- a/src/main/java/com/databricks/jdbc/api/internal/IDatabricksConnectionContext.java
+++ b/src/main/java/com/databricks/jdbc/api/internal/IDatabricksConnectionContext.java
@@ -426,6 +426,9 @@ default int getHeartbeatIntervalSeconds() {
/** Returns whether batched INSERT optimization is enabled */
boolean isBatchedInsertsEnabled();
+ /** Returns whether native parameter batch execution is enabled */
+ boolean isNativeBatchingEnabled();
+
/** Returns whether transaction-related method calls should be ignored */
boolean getIgnoreTransactions();
diff --git a/src/main/java/com/databricks/jdbc/common/DatabricksJdbcUrlParams.java b/src/main/java/com/databricks/jdbc/common/DatabricksJdbcUrlParams.java
index 7fea2fe2c..bf6531432 100644
--- a/src/main/java/com/databricks/jdbc/common/DatabricksJdbcUrlParams.java
+++ b/src/main/java/com/databricks/jdbc/common/DatabricksJdbcUrlParams.java
@@ -194,6 +194,7 @@ public enum DatabricksJdbcUrlParams {
"Timeout in seconds for metadata polling operations (e.g. GetTables, GetColumns). 0 means no timeout",
"300"),
ENABLE_BATCHED_INSERTS("EnableBatchedInserts", "Enable batched INSERT optimization", "0"),
+ ENABLE_NATIVE_BATCHING("EnableNativeBatching", "Enable native parameter batch execution", "0"),
ENABLE_SQL_VALIDATION_FOR_IS_VALID(
"EnableSQLValidationForIsValid",
"Enable SQL query execution for connection validation in isValid() method",
diff --git a/src/main/java/com/databricks/jdbc/common/util/ProtocolFeatureUtil.java b/src/main/java/com/databricks/jdbc/common/util/ProtocolFeatureUtil.java
index a451d7f0e..1143771fc 100644
--- a/src/main/java/com/databricks/jdbc/common/util/ProtocolFeatureUtil.java
+++ b/src/main/java/com/databricks/jdbc/common/util/ProtocolFeatureUtil.java
@@ -140,6 +140,16 @@ public static boolean supportsAsyncMetadataOperations(TProtocolVersion protocolV
return protocolVersion.compareTo(TProtocolVersion.SPARK_CLI_SERVICE_PROTOCOL_V9) >= 0;
}
+ /**
+ * Checks if the given protocol version supports native parameter batches.
+ *
+ * @param protocolVersion The protocol version to check
+ * @return true if native parameter batches are supported, false otherwise
+ */
+ public static boolean supportsNativeParameterBatching(TProtocolVersion protocolVersion) {
+ return protocolVersion.compareTo(TProtocolVersion.SPARK_CLI_SERVICE_PROTOCOL_V10) >= 0;
+ }
+
/**
* Checks if the given protocol version indicates a non-Databricks compute.
*
diff --git a/src/main/java/com/databricks/jdbc/dbclient/IDatabricksClient.java b/src/main/java/com/databricks/jdbc/dbclient/IDatabricksClient.java
index 71e790074..03b2083e3 100644
--- a/src/main/java/com/databricks/jdbc/dbclient/IDatabricksClient.java
+++ b/src/main/java/com/databricks/jdbc/dbclient/IDatabricksClient.java
@@ -15,6 +15,8 @@
import com.databricks.jdbc.telemetry.latency.DatabricksMetricsTimed;
import com.databricks.sdk.core.DatabricksConfig;
import java.sql.SQLException;
+import java.sql.SQLFeatureNotSupportedException;
+import java.util.List;
import java.util.Map;
/** Interface for Databricks client which abstracts the integration with Databricks server. */
@@ -71,6 +73,37 @@ DatabricksResultSet executeStatement(
MetadataOperationType metadataOperationType)
throws SQLException;
+ /**
+ * Returns whether this client can execute a native parameter batch for the given compute.
+ *
+ * @param computeResource underlying SQL warehouse or all-purpose cluster
+ */
+ default boolean supportsNativeParameterBatching(IDatabricksComputeResource computeResource) {
+ return false;
+ }
+
+ /**
+ * Executes one statement with multiple ordered parameter sets in a single backend request.
+ *
+ * @param sql SQL statement that needs to be executed
+ * @param computeResource underlying SQL warehouse or all-purpose cluster
+ * @param parameterSets ordered parameter sets for the statement
+ * @param statementType type of statement
+ * @param session underlying session
+ * @param parentStatement statement instance
+ */
+ @DatabricksMetricsTimed
+ default DatabricksResultSet executeStatementBatch(
+ String sql,
+ IDatabricksComputeResource computeResource,
+ List parameterSets,
+ StatementType statementType,
+ IDatabricksSession session,
+ IDatabricksStatementInternal parentStatement)
+ throws SQLException {
+ throw new SQLFeatureNotSupportedException("Native parameter batching is not supported");
+ }
+
/**
* Executes a statement in Databricks server asynchronously
*
diff --git a/src/main/java/com/databricks/jdbc/dbclient/impl/thrift/DatabricksThriftServiceClient.java b/src/main/java/com/databricks/jdbc/dbclient/impl/thrift/DatabricksThriftServiceClient.java
index 1eb7c82aa..9de6a3b34 100644
--- a/src/main/java/com/databricks/jdbc/dbclient/impl/thrift/DatabricksThriftServiceClient.java
+++ b/src/main/java/com/databricks/jdbc/dbclient/impl/thrift/DatabricksThriftServiceClient.java
@@ -13,9 +13,11 @@
import com.databricks.jdbc.api.internal.IDatabricksConnectionContext;
import com.databricks.jdbc.api.internal.IDatabricksSession;
import com.databricks.jdbc.api.internal.IDatabricksStatementInternal;
+import com.databricks.jdbc.common.AllPurposeCluster;
import com.databricks.jdbc.common.IDatabricksComputeResource;
import com.databricks.jdbc.common.MetadataOperationType;
import com.databricks.jdbc.common.StatementType;
+import com.databricks.jdbc.common.Warehouse;
import com.databricks.jdbc.common.util.DatabricksThreadContextHolder;
import com.databricks.jdbc.common.util.DriverUtil;
import com.databricks.jdbc.common.util.ProtocolFeatureUtil;
@@ -172,6 +174,33 @@ public DatabricksResultSet executeStatement(
return thriftAccessor.execute(request, parentStatement, session, statementType);
}
+ @Override
+ public boolean supportsNativeParameterBatching(IDatabricksComputeResource computeResource) {
+ if (computeResource instanceof AllPurposeCluster) {
+ return ProtocolFeatureUtil.supportsNativeParameterBatching(serverProtocolVersion);
+ }
+ return computeResource instanceof Warehouse;
+ }
+
+ @Override
+ public DatabricksResultSet executeStatementBatch(
+ String sql,
+ IDatabricksComputeResource computeResource,
+ List parameterSets,
+ StatementType statementType,
+ IDatabricksSession session,
+ IDatabricksStatementInternal parentStatement)
+ throws SQLException {
+ LOGGER.debug(
+ "Executing native parameter batch with {} parameter sets on {}",
+ parameterSets.size(),
+ computeResource);
+ DatabricksThreadContextHolder.setStatementType(statementType);
+ TExecuteStatementReq request =
+ getBatchRequest(sql, parameterSets, session, parentStatement, statementType);
+ return thriftAccessor.execute(request, parentStatement, session, statementType);
+ }
+
@Override
public DatabricksResultSet executeStatementAsync(
String sql,
@@ -194,17 +223,47 @@ public DatabricksResultSet executeStatementAsync(
@VisibleForTesting
TSparkParameter mapToSparkParameterListItem(ImmutableSqlParameter parameter) {
+ return mapToSparkParameterListItem(parameter, parameter.cardinal());
+ }
+
+ private TSparkParameter mapToSparkParameterListItem(
+ ImmutableSqlParameter parameter, int ordinal) {
Object value = parameter.value();
String typeString = parameter.type().name();
if (typeString.equals(DECIMAL) && value instanceof BigDecimal) {
typeString = getDecimalTypeString((BigDecimal) value);
}
return new TSparkParameter()
- .setOrdinal(parameter.cardinal())
+ .setOrdinal(ordinal)
.setType(typeString)
.setValue(value != null ? TSparkParameterValue.stringValue(value.toString()) : null);
}
+ private TExecuteStatementReq getBatchRequest(
+ String sql,
+ List parameterSets,
+ IDatabricksSession session,
+ IDatabricksStatementInternal parentStatement,
+ StatementType statementType)
+ throws SQLException {
+ TExecuteStatementReq request =
+ getRequest(sql, Collections.emptyMap(), session, parentStatement, false, statementType);
+ request.unsetParameters();
+ request.unsetResultRowLimit();
+ List> batchParameters =
+ parameterSets.stream()
+ .map(
+ parameterSet ->
+ parameterSet.getParameters().stream()
+ .map(
+ parameter ->
+ mapToSparkParameterListItem(parameter, parameter.cardinal() - 1))
+ .collect(Collectors.toList()))
+ .collect(Collectors.toList());
+ request.setBatchParameters(batchParameters);
+ return request;
+ }
+
private TExecuteStatementReq getRequest(
String sql,
Map parameters,
diff --git a/src/test/java/com/databricks/jdbc/api/impl/BatchParameterSetTest.java b/src/test/java/com/databricks/jdbc/api/impl/BatchParameterSetTest.java
new file mode 100644
index 000000000..dccb2e79c
--- /dev/null
+++ b/src/test/java/com/databricks/jdbc/api/impl/BatchParameterSetTest.java
@@ -0,0 +1,113 @@
+package com.databricks.jdbc.api.impl;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import com.databricks.jdbc.model.core.ColumnInfoTypeName;
+import java.sql.Timestamp;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import org.junit.jupiter.api.Test;
+
+class BatchParameterSetTest {
+
+ @Test
+ void ordersParametersAndPreservesJdbcIndexes() {
+ Map bindings = new HashMap<>();
+ bindings.put(3, parameter(99, "third", ColumnInfoTypeName.STRING));
+ bindings.put(1, parameter(99, "first", ColumnInfoTypeName.STRING));
+ bindings.put(2, parameter(99, "second", ColumnInfoTypeName.STRING));
+
+ BatchParameterSet parameterSet = BatchParameterSet.from(bindings);
+
+ assertEquals(List.of("first", "second", "third"), values(parameterSet));
+ assertEquals(List.of(1, 2, 3), indexes(parameterSet));
+ assertEquals(List.of(1, 2, 3), List.copyOf(parameterSet.getParameterBindings().keySet()));
+ }
+
+ @Test
+ void preservesSparseIndexesWithoutValidation() {
+ Map bindings = new HashMap<>();
+ bindings.put(3, parameter(3, "third", ColumnInfoTypeName.STRING));
+ bindings.put(1, parameter(1, "first", ColumnInfoTypeName.STRING));
+
+ BatchParameterSet parameterSet = BatchParameterSet.from(bindings);
+
+ assertEquals(List.of("first", "third"), values(parameterSet));
+ assertEquals(List.of(1, 3), indexes(parameterSet));
+ }
+
+ @Test
+ void allowsEmptyParameterSet() {
+ BatchParameterSet parameterSet = BatchParameterSet.from(Map.of());
+
+ assertTrue(parameterSet.isEmpty());
+ assertEquals(0, parameterSet.size());
+ }
+
+ @Test
+ void snapshotsBindingsAndMutableValues() {
+ Timestamp timestamp = Timestamp.valueOf("2026-08-10 12:34:56.123456789");
+ byte[] bytes = new byte[] {1, 2, 3};
+ Map bindings = new HashMap<>();
+ bindings.put(1, parameter(1, timestamp, ColumnInfoTypeName.TIMESTAMP));
+ bindings.put(2, parameter(2, bytes, ColumnInfoTypeName.BINARY));
+
+ BatchParameterSet parameterSet = BatchParameterSet.from(bindings);
+ bindings.clear();
+ timestamp.setTime(0);
+ bytes[0] = 9;
+
+ assertFalse(parameterSet.isEmpty());
+ assertEquals(
+ Timestamp.valueOf("2026-08-10 12:34:56.123456789"),
+ parameterSet.getParameters().get(0).value());
+ assertArrayEquals(new byte[] {1, 2, 3}, (byte[]) parameterSet.getParameters().get(1).value());
+ assertThrows(
+ UnsupportedOperationException.class,
+ () -> parameterSet.getParameters().add(parameter(3, "extra", ColumnInfoTypeName.STRING)));
+ assertThrows(
+ UnsupportedOperationException.class,
+ () ->
+ parameterSet
+ .getParameterBindings()
+ .put(3, parameter(3, "extra", ColumnInfoTypeName.STRING)));
+ }
+
+ @Test
+ void preservesNullValueAndType() {
+ BatchParameterSet parameterSet =
+ BatchParameterSet.from(Map.of(1, parameter(1, null, ColumnInfoTypeName.DECIMAL)));
+
+ ImmutableSqlParameter parameter = parameterSet.getParameters().get(0);
+ assertNull(parameter.value());
+ assertEquals(ColumnInfoTypeName.DECIMAL, parameter.type());
+ assertEquals(1, parameter.cardinal());
+ }
+
+ private ImmutableSqlParameter parameter(
+ int cardinal, Object value, ColumnInfoTypeName columnInfoTypeName) {
+ return ImmutableSqlParameter.builder()
+ .cardinal(cardinal)
+ .value(value)
+ .type(columnInfoTypeName)
+ .build();
+ }
+
+ private List