github-actions[bot] commented on code in PR #68768:
URL: https://github.com/apache/doris/pull/68768#discussion_r4215155737
##########
fe/fe-core/src/main/java/org/apache/doris/qe/ConnectContext.java:
##########
@@ -924,35 +931,157 @@ public void resetLoginTime() {
this.loginTime = System.currentTimeMillis();
}
- public synchronized void addPreparedQuery(String preparedStatementId,
String preparedQuery) {
- addPreparedQuery(preparedStatementId, preparedQuery, null);
+ public void addPreparedQuery(String preparedStatementId, String
preparedQuery) {
+ synchronized (preparedQuerys) {
+ addPreparedQuery(preparedStatementId, preparedQuery, null);
+ }
}
- public synchronized void addPreparedQuery(String preparedStatementId,
String preparedQuery, Schema schema) {
- preparedQuerys.put(preparedStatementId,
- new PreparedQuery(preparedQuery, getDefaultCatalog(),
getDatabase(), schema));
+ public void addPreparedQuery(String preparedStatementId, String
preparedQuery, Schema schema) {
+ synchronized (preparedQuerys) {
+ if (flightPreparedQueriesClosed) {
+ throw new IllegalStateException("Flight SQL session is
closed");
+ }
+ removePreparedQuery(preparedStatementId);
+ preparedQuerys.put(preparedStatementId,
+ new PreparedQuery(preparedQuery, getDefaultCatalog(),
getDatabase(), schema));
+ }
}
- public synchronized Schema getPreparedQuerySchema(String
preparedStatementId) {
- PreparedQuery query = preparedQuerys.get(preparedStatementId);
- return query == null ? null : query.schema;
+ public void addPreparedQuery(String id, String sql, Schema schema, int
parameterCount) {
+ synchronized (preparedQuerys) {
+ addPreparedQuery(id, sql, schema);
+ preparedQuerys.get(id).parameterCount = parameterCount;
+ }
}
- public synchronized String getPreparedQuery(String preparedStatementId) {
- PreparedQuery query = preparedQuerys.get(preparedStatementId);
- if (query == null) {
- return null;
+ public int getPreparedQueryParameterCount(String id) {
+ synchronized (preparedQuerys) {
+ PreparedQuery query = preparedQuerys.get(id);
+ return query == null ? -1 : query.parameterCount;
}
- // A handle must not execute unqualified SQL in a different namespace
than its advertised schema.
- if (!Objects.equals(query.catalog, getDefaultCatalog()) ||
!Objects.equals(query.database, getDatabase())) {
- preparedQuerys.remove(preparedStatementId);
- return null;
+ }
+
+ public List<Literal> getPreparedQueryParameters(String id) {
+ synchronized (preparedQuerys) {
+ PreparedQuery query = preparedQuerys.get(id);
+ return query == null ? null : query.parameters;
+ }
+ }
+
+ public long beginPreparedQueryBinding(String id) {
+ synchronized (preparedQuerys) {
+ PreparedQuery query = preparedQuerys.get(id);
+ if (query == null) {
+ return -1;
+ }
+ flightParameterBytes -= query.parameterBytes;
+ query.parameterBytes = 0;
+ query.parameters = null;
+ return ++query.bindingVersion;
}
- return query.sql;
}
- public synchronized void removePreparedQuery(String preparedStatementId) {
- preparedQuerys.remove(preparedStatementId);
+ public boolean isPreparedQueryBindingCurrent(String id, long version) {
+ synchronized (preparedQuerys) {
+ PreparedQuery query = preparedQuerys.get(id);
+ return query != null && query.bindingVersion == version;
+ }
+ }
+
+ public void setPreparedQueryParameters(String id, List<Literal>
parameters, Schema schema) {
+ synchronized (preparedQuerys) {
+ PreparedQuery query = preparedQuerys.get(id);
+ if (query == null) {
+ throw new IllegalStateException("Prepared statement expired");
+ }
+ long bytes = parameters == null ? 0 : parameters.size() * 128L;
+ if (parameters != null) {
+ for (Literal parameter : parameters) {
+ if (parameter.getValue() instanceof String) {
+ bytes += 2L * ((String) parameter.getValue()).length();
+ }
+ }
+ }
+ // A per-upload limit alone allows many handles to retain
unbounded parameter memory in one session.
+ if (flightParameterBytes - query.parameterBytes + bytes >
MAX_FLIGHT_PARAMETER_BYTES) {
+ throw new IllegalArgumentException("Prepared query parameters
exceed the 16 MiB session limit");
+ }
+ flightParameterBytes += bytes - query.parameterBytes;
+ query.parameterBytes = bytes;
+ // Retain detached immutable values, never vectors owned and
reused by the Flight stream.
+ query.parameters = parameters == null ? null :
List.copyOf(parameters);
+ query.schema = schema;
Review Comment:
[P2] Preserve the Prepare schema separately from the bound schema. For
`SELECT * FROM t WHERE id = ?`, a completed `ALTER TABLE t ADD COLUMN` between
Prepare and DoPut makes `analyzeQuerySchema` return a wider result, and this
assignment replaces the original baseline. `preparedQuery` then accepts a
handle whose result shape changed despite the existing schema-drift guard.
Unbound `GetSchema` also returns the changed shape without checking the
baseline. Reject incompatible DDL or session-setting drift in both paths while
allowing fields whose type genuinely depends on a bound value.
--
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]