Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
35 commits
Select commit Hold shift + click to select a range
85cc10b
fix: ensure CustomItemSerializer is not applied to internal query pip…
Apr 7, 2026
0a651aa
fix: address review iteration 1 initialize rawValue in ObjectNode co…
Apr 7, 2026
215ebf0
fix: clone SqlQuerySpec before applying serializer to prevent race co…
Apr 7, 2026
eb44dd6
fix: address PR review forceSerialization, remove ModelBridgeInterna…
Apr 7, 2026
0a01716
test: add coverage for CustomItemSerializer with ORDER BY, GROUP BY, …
Apr 7, 2026
01aa28d
fix: guard SqlParameter.applySerializer against unimplemented seriali…
Apr 7, 2026
6b4ddbe
fix: use clientSerializer instead of hardcoded EnvelopWrappingItemSer…
xinlian12 Apr 8, 2026
6f2f601
add changelog
xinlian12 Apr 8, 2026
0fa8c3a
fix: guard resolveTestNameSuffix against empty row for tests without …
xinlian12 Apr 8, 2026
2942745
fix: use DEFAULT_SERIALIZER for aggregate/distinct/groupby query tests
xinlian12 Apr 9, 2026
77f2b67
fix: use Integer.class for SELECT VALUE COUNT aggregate serializer test
xinlian12 Apr 10, 2026
dd17484
feat: add BasicCustomItemSerializer for query tests (issue #45521)
xinlian12 Apr 10, 2026
84b0de8
Skip query tests for envelope-wrapping serializer instead of falling …
xinlian12 Apr 10, 2026
fbc9e48
fix: CustomItemSerializer not applied correctly in queries and SqlPar…
xinlian12 Apr 13, 2026
f7124cc
fix: handle non-document values in EnvelopWrappingItemSerializer
xinlian12 Apr 13, 2026
2fa237a
fix: address review comments for custom serializer PR
xinlian12 Apr 13, 2026
7a26ae9
refactor: centralize serializer-neutralization for internal pipeline …
xinlian12 Apr 13, 2026
6835a4d
remove unused files
xinlian12 Apr 14, 2026
be6eb7a
update changelog
xinlian12 Apr 14, 2026
6cf01a1
Address PR review comments: NPE guard, optimization, Javadoc clarity
xinlian12 Apr 14, 2026
ad7c046
Fix Cosmos Spark BannedDependencies enforcer failure by updating scal…
xinlian12 Apr 14, 2026
7677a27
Merge branch 'fix/cosmos-spark-jackson-version-mismatch' into upstrea…
xinlian12 Apr 14, 2026
d205163
Merge branch 'main' of https://github.com/Azure/azure-sdk-for-java in…
xinlian12 Apr 14, 2026
73308ae
Merge branch 'main' of https://github.com/Azure/azure-sdk-for-java in…
xinlian12 Apr 14, 2026
63e649b
merge from main and resolve conflcits
xinlian12 Apr 14, 2026
0ffba39
Add concurrent SqlQuerySpec reuse test and improve pipeline Javadoc
xinlian12 Apr 14, 2026
0b15550
Merge branch 'main' of https://github.com/Azure/azure-sdk-for-java in…
xinlian12 Apr 15, 2026
b1d55b1
Merge branch 'upstream-main' into fix/issue-45521-custom-item-seriali…
xinlian12 Apr 15, 2026
0c7b6e5
fix: defensive copy of SqlQuerySpec in QueryPlanRetriever to prevent …
xinlian12 Apr 16, 2026
a021c31
refactor: add SqlQuerySpec.clone() and rewrite concurrent test with Flux
xinlian12 Apr 16, 2026
cbd5020
simplify: move SqlQuerySpec copy to private method in QueryPlanRetriever
xinlian12 Apr 16, 2026
cec9e1a
refactor: move cloneSqlParameter to SqlParameter.clone()
xinlian12 Apr 16, 2026
73c4360
refactor: make SqlParameter.createCopy() private with accessor pattern
xinlian12 Apr 16, 2026
a696f39
Merge branch 'main' of https://github.com/Azure/azure-sdk-for-java in…
xinlian12 Apr 17, 2026
c5339d3
Merge branch 'upstream-main' into fix/issue-45521-custom-item-seriali…
xinlian12 Apr 17, 2026
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
Original file line number Diff line number Diff line change
Expand Up @@ -9,4 +9,9 @@ private[cosmos] abstract class CosmosItemSerializerNoExceptionWrapping extends C
.CosmosItemSerializerHelper
.getCosmosItemSerializerAccessor
.setShouldWrapSerializationExceptions(this, false)

ImplementationBridgeHelpers
.CosmosItemSerializerHelper
.getCosmosItemSerializerAccessor
.setCanSerialize(this, false)
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,5 +14,10 @@ public CosmosItemSerializerNoExceptionWrapping() {
.CosmosItemSerializerHelper
.getCosmosItemSerializerAccessor()
.setShouldWrapSerializationExceptions(this, false);

ImplementationBridgeHelpers
.CosmosItemSerializerHelper
.getCosmosItemSerializerAccessor()
.setCanSerialize(this, false);
}
}

Large diffs are not rendered by default.

2 changes: 2 additions & 0 deletions sdk/cosmos/azure-cosmos/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@
* Fixed an issue where the throughput control `throughputQueryMono` was always subscribed even when `targetThroughput` is used (not `targetThroughputThreshold`), causing unnecessary `throughputSettings/read` permission requirement for AAD principals. - See [PR 48800](https://github.com/Azure/azure-sdk-for-java/pull/48800)
* Fixed an issue where change feed with `startFrom` point-in-time returned `400` on merged partitions by enabling the `CHANGE_FEED_WITH_START_TIME_POST_MERGE` SDK capability.
* Fixed JVM `<clinit>` deadlock when multiple threads concurrently trigger Cosmos SDK class loading for the first time. - See [PR 48689](https://github.com/Azure/azure-sdk-for-java/pull/48689)
* Fixed an issue where `CustomItemSerializer` was incorrectly applied to internal SDK query pipeline structures (e.g., `OrderByRowResult`, `Document`), causing deserialization failures in ORDER BY, GROUP BY, aggregate, DISTINCT, and hybrid search queries. - See [PR 48811](https://github.com/Azure/azure-sdk-for-java/pull/48811)
* Fixed an issue where `SqlParameter` ignored the configured `CustomItemSerializer`, always using the internal default serializer instead. - See [PR 48811](https://github.com/Azure/azure-sdk-for-java/pull/48811)
Comment thread
xinlian12 marked this conversation as resolved.

#### Other Changes

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,27 +3,16 @@

package com.azure.cosmos;

import com.azure.cosmos.implementation.ApiType;
import com.azure.cosmos.implementation.BadRequestException;
import com.azure.cosmos.implementation.Configs;
import com.azure.cosmos.implementation.ConnectionPolicy;
import com.azure.cosmos.implementation.CosmosClientMetadataCachesSnapshot;
import com.azure.cosmos.implementation.DefaultCosmosItemSerializer;
import com.azure.cosmos.implementation.HttpConstants;
import com.azure.cosmos.implementation.ImplementationBridgeHelpers;
import com.azure.cosmos.implementation.JsonSerializable;
import com.azure.cosmos.implementation.ObjectNodeMap;
import com.azure.cosmos.implementation.PrimitiveJsonNodeMap;
import com.azure.cosmos.implementation.Utils;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;

import java.io.IOException;
import java.util.Map;

import static com.azure.cosmos.implementation.guava25.base.Preconditions.checkNotNull;

/**
* The {@link CosmosItemSerializer} allows customizing the serialization of Cosmos Items - either to transform payload (for
* example wrap/unwrap in custom envelopes) or use custom serialization settings or json serializer stacks.
Expand All @@ -47,6 +36,7 @@ public abstract class CosmosItemSerializer {
new DefaultCosmosItemSerializer(Utils.getSimpleObjectMapper());

private boolean shouldWrapSerializationExceptions;
private boolean canSerialize;
Comment thread
xinlian12 marked this conversation as resolved.

private ObjectMapper mapper = Utils.getSimpleObjectMapper();

Expand All @@ -55,6 +45,7 @@ public abstract class CosmosItemSerializer {
*/
protected CosmosItemSerializer() {
this.shouldWrapSerializationExceptions = true;
this.canSerialize = true;
}

/**
Expand Down Expand Up @@ -139,6 +130,14 @@ void setShouldWrapSerializationExceptions(boolean enabled) {
this.shouldWrapSerializationExceptions = enabled;
}

void setCanSerialize(boolean canSerialize) {
this.canSerialize = canSerialize;
}

boolean canSerialize() {
return this.canSerialize;
}


///////////////////////////////////////////////////////////////////////////////////////////
// the following helper/accessor only helps to access this class outside of this package.//
Expand Down Expand Up @@ -172,6 +171,15 @@ public ObjectMapper getItemObjectMapper(CosmosItemSerializer serializer) {
}

@Override
public boolean canSerialize(CosmosItemSerializer serializer) {
return serializer.canSerialize();
}

@Override
public void setCanSerialize(CosmosItemSerializer serializer, boolean canSerialize) {
serializer.setCanSerialize(canSerialize);
}

public CosmosItemSerializer getInternalDefaultSerializer() {
return INTERNAL_DEFAULT_SERIALIZER;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,7 @@
import com.azure.cosmos.models.PartitionKeyDefinition;
import com.azure.cosmos.models.PriorityLevel;
import com.azure.cosmos.models.ShowQueryMode;
import com.azure.cosmos.models.SqlParameter;
import com.azure.cosmos.models.SqlQuerySpec;
import com.azure.cosmos.util.CosmosPagedFlux;
import com.azure.cosmos.util.UtilBridgeInternal;
Expand Down Expand Up @@ -1893,6 +1894,8 @@ void setShouldWrapSerializationExceptions(
void setItemObjectMapper(CosmosItemSerializer serializer, ObjectMapper mapper);
ObjectMapper getItemObjectMapper(CosmosItemSerializer serializer);
CosmosItemSerializer getInternalDefaultSerializer();
void setCanSerialize(CosmosItemSerializer serializer, boolean canSerialize);
boolean canSerialize(CosmosItemSerializer serializer);
}
}

Expand Down Expand Up @@ -1934,4 +1937,72 @@ ReadConsistencyStrategy getEffectiveReadConsistencyStrategy(
ReadConsistencyStrategy clientLevelReadConsistencyStrategy);
}
}

public static final class SqlQuerySpecHelper {
private static final AtomicReference<SqlQuerySpecAccessor> accessor = new AtomicReference<>();
private static final AtomicBoolean sqlQuerySpecClassLoaded = new AtomicBoolean(false);

private SqlQuerySpecHelper() {}

public static void setSqlQuerySpecAccessor(final SqlQuerySpecAccessor newAccessor) {
if (!accessor.compareAndSet(null, newAccessor)) {
logger.debug("SqlQuerySpecAccessor already initialized!");
} else {
logger.debug("Setting SqlQuerySpecAccessor...");
sqlQuerySpecClassLoaded.set(true);
}
}

public static SqlQuerySpecAccessor getSqlQuerySpecAccessor() {
if (!sqlQuerySpecClassLoaded.get()) {
logger.debug("Initializing SqlQuerySpecAccessor...");
initializeAllAccessors();
}

SqlQuerySpecAccessor snapshot = accessor.get();
if (snapshot == null) {
logger.error("SqlQuerySpecAccessor is not initialized yet!");
}

return snapshot;
}

public interface SqlQuerySpecAccessor {
void applySerializerToParameters(SqlQuerySpec sqlQuerySpec, CosmosItemSerializer serializer);
}
}

public static final class SqlParameterHelper {
private static final AtomicReference<SqlParameterAccessor> accessor = new AtomicReference<>();
private static final AtomicBoolean sqlParameterClassLoaded = new AtomicBoolean(false);

private SqlParameterHelper() {}

public static void setSqlParameterAccessor(final SqlParameterAccessor newAccessor) {
if (!accessor.compareAndSet(null, newAccessor)) {
logger.debug("SqlParameterAccessor already initialized!");
} else {
logger.debug("Setting SqlParameterAccessor...");
sqlParameterClassLoaded.set(true);
}
}

public static SqlParameterAccessor getSqlParameterAccessor() {
if (!sqlParameterClassLoaded.get()) {
logger.debug("Initializing SqlParameterAccessor...");
initializeAllAccessors();
}

SqlParameterAccessor snapshot = accessor.get();
if (snapshot == null) {
logger.error("SqlParameterAccessor is not initialized yet!");
}

return snapshot;
}

public interface SqlParameterAccessor {
SqlParameter createCopy(SqlParameter original);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -63,29 +63,26 @@ private static BiFunction<String, PipelinedDocumentQueryParams<Document>, Flux<I
if (queryInfo.hasOrderBy()) {
createBaseComponentFunction = (continuationToken, documentQueryParams) -> {
CosmosQueryRequestOptions orderByCosmosQueryRequestOptions =
qryOptAccessor().clone(requestOptions);
cloneOptionsForInternalPipeline(requestOptions);
if (queryInfo.hasNonStreamingOrderBy()) {
if (continuationToken != null) {
throw new NonStreamingOrderByBadRequestException(
HttpConstants.StatusCodes.BADREQUEST,
"Can not use a continuation token for a vector search query");
}
qryOptAccessor().getImpl(orderByCosmosQueryRequestOptions).setCustomItemSerializer(null);

documentQueryParams.setCosmosQueryRequestOptions(orderByCosmosQueryRequestOptions);
return NonStreamingOrderByDocumentQueryExecutionContext.createAsync(diagnosticsClientContext, client, documentQueryParams, collection);
} else {
ModelBridgeInternal.setQueryRequestOptionsContinuationToken(orderByCosmosQueryRequestOptions, continuationToken);
qryOptAccessor().getImpl(orderByCosmosQueryRequestOptions).setCustomItemSerializer(null);
documentQueryParams.setCosmosQueryRequestOptions(orderByCosmosQueryRequestOptions);
return OrderByDocumentQueryExecutionContext.createAsync(diagnosticsClientContext, client, documentQueryParams, collection);
}
};
} else {

createBaseComponentFunction = (continuationToken, documentQueryParams) -> {
CosmosQueryRequestOptions parallelCosmosQueryRequestOptions =
qryOptAccessor().clone(requestOptions);
qryOptAccessor().getImpl(parallelCosmosQueryRequestOptions).setCustomItemSerializer(null);
CosmosQueryRequestOptions parallelCosmosQueryRequestOptions = cloneOptionsForInternalPipeline(requestOptions);
ModelBridgeInternal.setQueryRequestOptionsContinuationToken(parallelCosmosQueryRequestOptions, continuationToken);

documentQueryParams.setCosmosQueryRequestOptions(parallelCosmosQueryRequestOptions);
Expand Down Expand Up @@ -122,7 +119,7 @@ private static BiFunction<String, PipelinedDocumentQueryParams<Document>, Flux<I

createBaseComponentFunction = (continuationToken, documentQueryParams) -> {
CosmosQueryRequestOptions orderByCosmosQueryRequestOptions =
qryOptAccessor().clone(requestOptions);
cloneOptionsForInternalPipeline(requestOptions);

documentQueryParams.setCosmosQueryRequestOptions(orderByCosmosQueryRequestOptions);
return HybridSearchDocumentQueryExecutionContext.createAsync(diagnosticsClientContext, client, documentQueryParams, collection);
Expand Down Expand Up @@ -273,4 +270,24 @@ public Flux<FeedResponse<T>> executeAsync() {
)
);
}

/**
* Clones query request options and neutralizes the custom item serializer to
* DEFAULT_SERIALIZER. Internal pipeline stages must always use the default
* serializer; custom serialization is applied at the top-level pipeline boundary.
*
* All internal pipeline stages that work with Document/intermediate types MUST
* use this method rather than cloning options directly via
* {@code qryOptAccessor().clone(...)}. Bypassing this method would leave the
* custom serializer active in the inner pipeline, causing incorrect
* serialization of intermediate results. See PR #48811 for context.
*/
private static CosmosQueryRequestOptions cloneOptionsForInternalPipeline(
Comment thread
xinlian12 marked this conversation as resolved.
CosmosQueryRequestOptions source) {

CosmosQueryRequestOptions cloned = qryOptAccessor().clone(source);
qryOptAccessor().getImpl(cloned)
.setCustomItemSerializer(CosmosItemSerializer.DEFAULT_SERIALIZER);
return cloned;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ public class PipelinedDocumentQueryParams<T> {
private final String collectionRid;
private final ResourceType resourceTypeEnum;
private final Class<T> resourceType;
private final SqlQuerySpec query;
private SqlQuerySpec query;
Comment thread
xinlian12 marked this conversation as resolved.
private final String resourceLink;
private final UUID correlatedActivityId;
private CosmosQueryRequestOptions cosmosQueryRequestOptions;
Expand Down Expand Up @@ -101,6 +101,10 @@ public SqlQuerySpec getQuery() {
return query;
}

public void setQuery(SqlQuerySpec query) {
Comment thread
xinlian12 marked this conversation as resolved.
this.query = query;
}

public String getResourceLink() {
return resourceLink;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,12 +5,20 @@
import com.azure.cosmos.CosmosItemSerializer;
import com.azure.cosmos.implementation.DiagnosticsClientContext;
import com.azure.cosmos.implementation.DocumentCollection;
import com.azure.cosmos.implementation.ImplementationBridgeHelpers;
import com.azure.cosmos.implementation.ImplementationBridgeHelpers.CosmosItemSerializerHelper.CosmosItemSerializerAccessor;
import com.azure.cosmos.implementation.ImplementationBridgeHelpers.SqlParameterHelper.SqlParameterAccessor;
import com.azure.cosmos.implementation.ImplementationBridgeHelpers.SqlQuerySpecHelper.SqlQuerySpecAccessor;
import com.azure.cosmos.implementation.Utils;
import com.azure.cosmos.implementation.query.hybridsearch.HybridSearchQueryInfo;
import com.azure.cosmos.models.CosmosQueryRequestOptions;
import com.azure.cosmos.models.ModelBridgeInternal;
import com.azure.cosmos.models.SqlParameter;
import com.azure.cosmos.models.SqlQuerySpec;
import reactor.core.publisher.Flux;

import java.util.ArrayList;
import java.util.List;
import java.util.function.BiFunction;

/**
Expand Down Expand Up @@ -60,6 +68,28 @@ public static <T> Flux<PipelinedQueryExecutionContextBase<T>> createAsync(

CosmosItemSerializer candidateSerializer = client.getEffectiveItemSerializer(cosmosQueryRequestOptions);

// Apply the effective custom serializer to SqlParameter values so that
// query parameters are serialized consistently with stored document values.
// Clone the SqlQuerySpec first to avoid mutating the caller's original object,
// which could race if the same instance is reused across concurrent queries.
Comment thread
xinlian12 marked this conversation as resolved.
if (candidateSerializer != CosmosItemSerializer.DEFAULT_SERIALIZER
&& cosmosItemSerializerAccessor().canSerialize(candidateSerializer)) {
SqlQuerySpec original = initParams.getQuery();
if (original != null) {
List<SqlParameter> originalParams = original.getParameters();
if (originalParams != null && !originalParams.isEmpty()) {
List<SqlParameter> clonedParams = new ArrayList<>(originalParams.size());
SqlParameterAccessor paramAccessor = sqlParameterAccessor();
for (SqlParameter p : originalParams) {
clonedParams.add(paramAccessor.createCopy(p));
}
SqlQuerySpec clonedQuery = new SqlQuerySpec(original.getQueryText(), clonedParams);
sqlQuerySpecAccessor().applySerializerToParameters(clonedQuery, candidateSerializer);
initParams.setQuery(clonedQuery);
}
}
}

if (hybridSearchQueryInfo != null) {
return PipelinedDocumentQueryExecutionContext.createHybridAsyncCore(
diagnosticsClientContext,
Expand Down Expand Up @@ -194,4 +224,16 @@ public QueryInfo getQueryInfo() {
public HybridSearchQueryInfo getHybridSearchQueryInfo() {
return this.hybridSearchQueryInfo;
}

private static SqlQuerySpecAccessor sqlQuerySpecAccessor() {
return ImplementationBridgeHelpers.SqlQuerySpecHelper.getSqlQuerySpecAccessor();
}

private static SqlParameterAccessor sqlParameterAccessor() {
return ImplementationBridgeHelpers.SqlParameterHelper.getSqlParameterAccessor();
}

private static CosmosItemSerializerAccessor cosmosItemSerializerAccessor() {
return ImplementationBridgeHelpers.CosmosItemSerializerHelper.getCosmosItemSerializerAccessor();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
import com.azure.cosmos.models.CosmosQueryRequestOptions;
import com.azure.cosmos.models.ModelBridgeInternal;
import com.azure.cosmos.models.PartitionKey;
import com.azure.cosmos.models.SqlParameter;
import com.azure.cosmos.models.SqlQuerySpec;
import com.azure.cosmos.implementation.BackoffRetryUtility;
import com.azure.cosmos.implementation.DocumentClientRetryPolicy;
Expand All @@ -26,6 +27,7 @@
import com.fasterxml.jackson.databind.node.ObjectNode;
import reactor.core.publisher.Mono;

import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
Expand Down Expand Up @@ -105,7 +107,13 @@ static Mono<PartitionedQueryExecutionInfo> getQueryPlanThroughGatewayAsync(Diagn
resourceLink,
requestHeaders);
queryPlanRequest.useGatewayMode = true;
queryPlanRequest.setByteBuffer(ModelBridgeInternal.serializeJsonToByteBuffer(sqlQuerySpec));

// Create a defensive copy to prevent concurrent modification of the shared
// SqlQuerySpec's internal ObjectNode when multiple threads retrieve the query
// plan simultaneously. Each copy has its own property bag, avoiding the race
// condition on the non-thread-safe ObjectNode/LinkedHashMap backing store.
SqlQuerySpec querySpecCopy = copyQuerySpec(sqlQuerySpec);
queryPlanRequest.setByteBuffer(ModelBridgeInternal.serializeJsonToByteBuffer(querySpecCopy));

CosmosEndToEndOperationLatencyPolicyConfig end2EndConfig = queryOptionsAccessor()
.getImpl(nonNullRequestOptions)
Expand Down Expand Up @@ -162,4 +170,10 @@ static Mono<PartitionedQueryExecutionInfo> getQueryPlanThroughGatewayAsync(Diagn
return throwable;
});
}

private static SqlQuerySpec copyQuerySpec(SqlQuerySpec original) {
List<SqlParameter> params = original.getParameters();
List<SqlParameter> copiedParams = params != null ? new ArrayList<>(params) : null;
return new SqlQuerySpec(original.getQueryText(), copiedParams);
}
}
Loading
Loading