Skip to content
This repository was archived by the owner on Jul 13, 2026. It is now read-only.

Commit ec69803

Browse files
committed
add sync sample and other updates
1 parent 0835ad6 commit ec69803

5 files changed

Lines changed: 274 additions & 126 deletions

File tree

src/main/java/com/azure/cosmos/examples/bulk/async/SampleBulkHandleRetriesAsync.java

Lines changed: 2 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -86,7 +86,8 @@ private void getStartedDemo() {
8686
.consistencyLevel(ConsistencyLevel.SESSION).buildAsyncClient();
8787

8888

89-
createDatabaseIfNotExists();
89+
// Note: database must already exist (cannot be created via RBAC).
90+
database = client.getDatabase(databaseName);
9091
createContainerIfNotExists();
9192

9293
// Create 500 Family records
@@ -102,20 +103,6 @@ private void getStartedDemo() {
102103
largeBulkUpsertItemsWithBulkWriterAbstraction(largeFamilies1);
103104
}
104105

105-
private void createDatabaseIfNotExists() {
106-
logger.info("Create database " + databaseName + " if not exists.");
107-
108-
// Create database if not exists
109-
// <CreateDatabaseIfNotExists>
110-
Mono<CosmosDatabaseResponse> databaseIfNotExists = client.createDatabaseIfNotExists(databaseName);
111-
databaseIfNotExists.flatMap(databaseResponse -> {
112-
database = client.getDatabase(databaseResponse.getProperties().getId());
113-
logger.info("Checking database " + database.getId() + " completed!\n");
114-
return Mono.empty();
115-
}).block();
116-
// </CreateDatabaseIfNotExists>
117-
}
118-
119106
private void createContainerIfNotExists() {
120107
logger.info("Create container " + containerName + " if not exists.");
121108

src/main/java/com/azure/cosmos/examples/bulk/async/SampleBulkQuickStartAsync.java

Lines changed: 4 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@
1010
// <CosmosBulkOperationsImport>
1111
import com.azure.cosmos.models.*;
1212
// </CosmosBulkOperationsImport>
13+
import com.azure.identity.DefaultAzureCredentialBuilder;
1314
import org.slf4j.Logger;
1415
import org.slf4j.LoggerFactory;
1516
import reactor.core.publisher.Flux;
@@ -63,14 +64,15 @@ private void getStartedDemo() {
6364
// <CreateAsyncClient>
6465
client = new CosmosClientBuilder()
6566
.endpoint(AccountSettings.HOST)
66-
.key(AccountSettings.MASTER_KEY)
67+
.credential(new DefaultAzureCredentialBuilder().build())
6768
.preferredRegions(preferredRegions)
6869
.contentResponseOnWriteEnabled(true)
6970
.consistencyLevel(ConsistencyLevel.SESSION).buildAsyncClient();
7071

7172
// </CreateAsyncClient>
7273

73-
createDatabaseIfNotExists();
74+
// Note: database must already exist (cannot be created via RBAC).
75+
database = client.getDatabase(databaseName);
7476
createContainerIfNotExists();
7577

7678
// <AddDocsToStream>
@@ -131,20 +133,6 @@ private void getStartedDemo() {
131133
bulkCreateItemsWithBulkWriterAbstractionAndGlobalThroughputControl();
132134
}
133135

134-
private void createDatabaseIfNotExists() {
135-
logger.info("Create database " + databaseName + " if not exists.");
136-
137-
// Create database if not exists
138-
// <CreateDatabaseIfNotExists>
139-
Mono<CosmosDatabaseResponse> databaseIfNotExists = client.createDatabaseIfNotExists(databaseName);
140-
databaseIfNotExists.flatMap(databaseResponse -> {
141-
database = client.getDatabase(databaseResponse.getProperties().getId());
142-
logger.info("Checking database " + database.getId() + " completed!\n");
143-
return Mono.empty();
144-
}).block();
145-
// </CreateDatabaseIfNotExists>
146-
}
147-
148136
private void createContainerIfNotExists() {
149137
logger.info("Create container " + containerName + " if not exists.");
150138

@@ -349,8 +337,6 @@ private void shutdown() {
349337
logger.info("Deleting Cosmos DB resources");
350338
logger.info("-Deleting container...");
351339
if (container != null) container.delete().subscribe();
352-
logger.info("-Deleting database...");
353-
if (database != null) database.delete().subscribe();
354340
logger.info("-Closing the client...");
355341
} catch (InterruptedException err) {
356342
err.printStackTrace();
Lines changed: 97 additions & 78 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,4 @@
1-
// Copyright (c) Microsoft Corporation. All rights reserved.
2-
// Licensed under the MIT License.
31

4-
/*
5-
The BulkWriter class is an attempt to provide guidance for creating
6-
a higher level abstraction over the existing low level Java Bulk API
7-
*/
82
package com.azure.cosmos.examples.bulk.sync;
93

104
import com.azure.cosmos.CosmosContainer;
@@ -16,134 +10,159 @@
1610
import com.azure.cosmos.models.CosmosItemOperation;
1711
import org.slf4j.Logger;
1812
import org.slf4j.LoggerFactory;
19-
import reactor.core.publisher.Sinks;
13+
14+
import java.util.ArrayList;
15+
import java.util.List;
2016
import java.util.concurrent.Semaphore;
2117

2218
public class BulkWriter {
2319
private static final Logger logger = LoggerFactory.getLogger(BulkWriter.class);
2420

25-
private final Sinks.Many<CosmosItemOperation> bulkInputEmitter = Sinks.many().unicast().onBackpressureBuffer();
21+
private final CosmosContainer cosmosContainer;
2622
private final int cpuCount = Runtime.getRuntime().availableProcessors();
27-
28-
//Max items to be buffered to avoid out of memory error
2923
private final Semaphore semaphore = new Semaphore(1024 * 167 / cpuCount);
30-
31-
private final Sinks.EmitFailureHandler emitFailureHandler =
32-
(signalType, emitResult) -> {
33-
if (emitResult.equals(Sinks.EmitResult.FAIL_NON_SERIALIZED)) {
34-
logger.debug("emitFailureHandler - Signal: [{}], Result: [{}]", signalType, emitResult);
35-
return true;
36-
} else {
37-
logger.error("emitFailureHandler - Signal: [{}], Result: [{}]", signalType, emitResult);
38-
return false;
39-
}
40-
};
41-
42-
private final CosmosContainer cosmosContainer;
24+
private final List<CosmosItemOperation> bufferedOperations = new ArrayList<>();
25+
private final int maxRetries = 5;
4326

4427
public BulkWriter(CosmosContainer cosmosContainer) {
4528
this.cosmosContainer = cosmosContainer;
4629
}
4730

4831
public void scheduleWrites(CosmosItemOperation cosmosItemOperation) {
49-
while(!semaphore.tryAcquire()) {
32+
while (!semaphore.tryAcquire()) {
5033
logger.info("Unable to acquire permit");
5134
}
5235
logger.info("Acquired permit");
53-
scheduleInternalWrites(cosmosItemOperation);
54-
}
55-
56-
private void scheduleInternalWrites(CosmosItemOperation cosmosItemOperation) {
57-
bulkInputEmitter.emitNext(cosmosItemOperation, emitFailureHandler);
36+
bufferedOperations.add(cosmosItemOperation);
5837
}
5938

6039
public Iterable<CosmosBulkOperationResponse<Object>> execute() {
61-
return this.execute(null);
40+
return execute(null);
6241
}
6342

6443
public Iterable<CosmosBulkOperationResponse<Object>> execute(CosmosBulkExecutionOptions bulkOptions) {
6544
if (bulkOptions == null) {
6645
bulkOptions = new CosmosBulkExecutionOptions();
6746
}
68-
bulkInputEmitter.tryEmitComplete();
69-
Iterable<CosmosBulkOperationResponse<Object>> bulkOperationResponse = cosmosContainer
70-
.executeBulkOperations(
71-
bulkInputEmitter.asFlux().toIterable(),
72-
bulkOptions);
73-
for (CosmosBulkOperationResponse<Object> response : bulkOperationResponse) {
74-
processBulkOperationResponse(
75-
response.getResponse(),
76-
response.getOperation(),
77-
response.getException());
78-
}
79-
semaphore.release();
80-
return bulkOperationResponse;
81-
}
8247

48+
List<CosmosItemOperation> currentBatch = new ArrayList<>(bufferedOperations);
49+
bufferedOperations.clear();
50+
List<CosmosBulkOperationResponse<Object>> finalResponses = new ArrayList<>();
51+
int attempt = 0;
52+
53+
while (!currentBatch.isEmpty() && attempt <= maxRetries) {
54+
logger.info("Executing bulk attempt {} with {} items", attempt + 1, currentBatch.size());
55+
56+
Iterable<CosmosBulkOperationResponse<Object>> responses =
57+
cosmosContainer.executeBulkOperations(currentBatch, bulkOptions);
58+
59+
List<CosmosItemOperation> toRetry = new ArrayList<>();
60+
61+
for (CosmosBulkOperationResponse<Object> response : responses) {
62+
processBulkOperationResponse(
63+
response.getResponse(),
64+
response.getOperation(),
65+
response.getException(),
66+
toRetry
67+
);
68+
finalResponses.add(response);
69+
}
70+
71+
currentBatch = toRetry;
72+
attempt++;
73+
}
8374

75+
if (!currentBatch.isEmpty()) {
76+
logger.error("BulkWriter: {} items failed after {} retries", currentBatch.size(), maxRetries);
77+
}
8478

79+
semaphore.release();
80+
return finalResponses;
81+
}
8582

8683
private void processBulkOperationResponse(
87-
CosmosBulkItemResponse itemResponse,
88-
CosmosItemOperation itemOperation,
89-
Exception exception) {
84+
CosmosBulkItemResponse itemResponse,
85+
CosmosItemOperation itemOperation,
86+
Exception exception,
87+
List<CosmosItemOperation> retryList) {
9088

9189
if (exception != null) {
92-
handleException(itemOperation, exception);
90+
handleException(itemOperation, exception, retryList);
9391
} else {
94-
processResponseCode(itemResponse, itemOperation);
92+
processResponseCode(itemResponse, itemOperation, retryList);
9593
}
9694
}
9795

9896
private void processResponseCode(
99-
CosmosBulkItemResponse itemResponse,
100-
CosmosItemOperation itemOperation) {
97+
CosmosBulkItemResponse itemResponse,
98+
CosmosItemOperation itemOperation,
99+
List<CosmosItemOperation> retryList) {
101100

102101
if (itemResponse.isSuccessStatusCode()) {
103102
logger.info(
104-
"The operation for Item ID: [{}] Item PartitionKey Value: [{}] completed successfully " +
105-
"with a response status code: [{}]",
106-
itemOperation.getId(),
107-
itemOperation.getPartitionKeyValue(),
108-
itemResponse.getStatusCode());
103+
"Item ID [{}] with PartitionKey [{}] succeeded with status [{}]",
104+
itemOperation.getId(),
105+
itemOperation.getPartitionKeyValue(),
106+
itemResponse.getStatusCode());
109107
} else if (shouldRetry(itemResponse.getStatusCode())) {
110108
logger.info(
111-
"The operation for Item ID: [{}] Item PartitionKey Value: [{}] will be retried",
112-
itemOperation.getId(),
113-
itemOperation.getPartitionKeyValue());
114-
//re-scheduling
115-
scheduleWrites(itemOperation);
109+
"Item ID [{}] with PartitionKey [{}] will be retried (status [{}])",
110+
itemOperation.getId(),
111+
itemOperation.getPartitionKeyValue(),
112+
itemResponse.getStatusCode());
113+
if (itemResponse.getRetryAfterDuration() != null) {
114+
logger.info(
115+
"Item ID [{}] with PartitionKey [{}] will be retried after [{}] milliseconds",
116+
itemOperation.getId(),
117+
itemOperation.getPartitionKeyValue(),
118+
itemResponse.getRetryAfterDuration().toMillis());
119+
try {
120+
Thread.sleep(itemResponse.getRetryAfterDuration().toMillis());
121+
} catch (InterruptedException e) {
122+
throw new RuntimeException(e);
123+
}
124+
}
125+
retryList.add(itemOperation);
116126
} else {
117127
logger.info(
118-
"The operation for Item ID: [{}] Item PartitionKey Value: [{}] did not complete successfully " +
119-
"with a response status code: [{}]",
120-
itemOperation.getId(),
121-
itemOperation.getPartitionKeyValue(),
122-
itemResponse.getStatusCode());
128+
"Item ID [{}] with PartitionKey [{}] failed with non-retryable status [{}]",
129+
itemOperation.getId(),
130+
itemOperation.getPartitionKeyValue(),
131+
itemResponse.getStatusCode());
123132
}
124133
}
125134

126-
private void handleException(CosmosItemOperation itemOperation, Exception exception) {
135+
private void handleException(
136+
CosmosItemOperation itemOperation,
137+
Exception exception,
138+
List<CosmosItemOperation> retryList) {
139+
127140
if (!(exception instanceof CosmosException)) {
128141
logger.info(
129-
"The operation for Item ID: [{}] Item PartitionKey Value: [{}] encountered an unexpected failure",
130-
itemOperation.getId(),
131-
itemOperation.getPartitionKeyValue());
132-
} else {
133-
if (shouldRetry(((CosmosException) exception).getStatusCode())) {
134-
logger.info(
135-
"The operation for Item ID: [{}] Item PartitionKey Value: [{}] will be retried",
142+
"Item ID [{}] with PartitionKey [{}] encountered unexpected failure",
136143
itemOperation.getId(),
137144
itemOperation.getPartitionKeyValue());
138-
139-
//re-scheduling
140-
scheduleWrites(itemOperation);
145+
} else {
146+
int statusCode = ((CosmosException) exception).getStatusCode();
147+
if (shouldRetry(statusCode)) {
148+
logger.info(
149+
"Item ID [{}] with PartitionKey [{}] will be retried due to exception status [{}]",
150+
itemOperation.getId(),
151+
itemOperation.getPartitionKeyValue(),
152+
statusCode);
153+
retryList.add(itemOperation);
154+
} else {
155+
logger.error(
156+
"Item ID [{}] with PartitionKey [{}] failed with non-retryable exception: {}",
157+
itemOperation.getId(),
158+
itemOperation.getPartitionKeyValue(),
159+
exception.getMessage());
141160
}
142161
}
143162
}
144163

145164
private boolean shouldRetry(int statusCode) {
146165
return statusCode == HttpConstants.StatusCodes.REQUEST_TIMEOUT ||
147-
statusCode == HttpConstants.StatusCodes.TOO_MANY_REQUESTS;
166+
statusCode == HttpConstants.StatusCodes.TOO_MANY_REQUESTS;
148167
}
149168
}

0 commit comments

Comments
 (0)