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

Commit 0835ad6

Browse files
committed
async handle retry sample
1 parent 59d21b9 commit 0835ad6

3 files changed

Lines changed: 222 additions & 0 deletions

File tree

pom.xml

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,11 @@
4747
<artifactId>azure-cosmos</artifactId>
4848
<version>LATEST</version>
4949
</dependency>
50+
<dependency>
51+
<groupId>com.azure</groupId>
52+
<artifactId>azure-identity</artifactId>
53+
<version>1.15.4</version>
54+
</dependency>
5055
<dependency>
5156
<groupId>org.apache.logging.log4j</groupId>
5257
<artifactId>log4j-api</artifactId>

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

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,19 @@ private void processResponseCode(
111111
"The operation for Item ID: [{}] Item PartitionKey Value: [{}] will be retried",
112112
itemOperation.getId(),
113113
itemOperation.getPartitionKeyValue());
114+
if (itemResponse.getRetryAfterDuration() != null) {
115+
logger.info(
116+
"Retrying after [{}] milliseconds",
117+
itemResponse.getRetryAfterDuration().toMillis());
118+
try {
119+
Thread.sleep(itemResponse.getRetryAfterDuration().toMillis());
120+
} catch (InterruptedException e) {
121+
throw new RuntimeException(e);
122+
}
123+
} else {
124+
logger.info("Retrying without delay");
125+
}
126+
114127
//re-scheduling
115128
scheduleWrites(itemOperation);
116129
} else {
Lines changed: 204 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,204 @@
1+
// Copyright (c) Microsoft Corporation. All rights reserved.
2+
// Licensed under the MIT License.
3+
4+
package com.azure.cosmos.examples.bulk.async;
5+
6+
import com.azure.cosmos.ConsistencyLevel;
7+
import com.azure.cosmos.CosmosAsyncClient;
8+
import com.azure.cosmos.CosmosAsyncContainer;
9+
import com.azure.cosmos.CosmosAsyncDatabase;
10+
import com.azure.cosmos.CosmosClientBuilder;
11+
import com.azure.cosmos.GlobalThroughputControlConfig;
12+
import com.azure.cosmos.ThrottlingRetryOptions;
13+
import com.azure.cosmos.ThroughputControlGroupConfig;
14+
import com.azure.cosmos.ThroughputControlGroupConfigBuilder;
15+
import com.azure.cosmos.examples.common.AccountSettings;
16+
import com.azure.cosmos.examples.common.Families;
17+
import com.azure.cosmos.examples.common.Family;
18+
import com.azure.cosmos.models.CosmosBulkExecutionOptions;
19+
import com.azure.cosmos.models.CosmosBulkItemResponse;
20+
import com.azure.cosmos.models.CosmosBulkOperations;
21+
import com.azure.cosmos.models.CosmosContainerProperties;
22+
import com.azure.cosmos.models.CosmosContainerRequestOptions;
23+
import com.azure.cosmos.models.CosmosContainerResponse;
24+
import com.azure.cosmos.models.CosmosDatabaseResponse;
25+
import com.azure.cosmos.models.CosmosItemOperation;
26+
import com.azure.cosmos.models.CosmosItemRequestOptions;
27+
import com.azure.cosmos.models.CosmosPatchOperations;
28+
import com.azure.cosmos.models.PartitionKey;
29+
import com.azure.cosmos.models.ThroughputProperties;
30+
import com.azure.identity.DefaultAzureCredentialBuilder;
31+
import org.slf4j.Logger;
32+
import org.slf4j.LoggerFactory;
33+
import reactor.core.publisher.Flux;
34+
import reactor.core.publisher.Mono;
35+
36+
import java.time.Duration;
37+
import java.util.ArrayList;
38+
import java.util.List;
39+
import java.util.UUID;
40+
41+
public class SampleBulkHandleRetriesAsync {
42+
43+
private static final Logger logger = LoggerFactory.getLogger(SampleBulkHandleRetriesAsync.class);
44+
private final String databaseName = "AzureSampleFamilyDB";
45+
private final String containerName = "FamilyContainer";
46+
private CosmosAsyncClient client;
47+
private CosmosAsyncDatabase database;
48+
private CosmosAsyncContainer container;
49+
private int noOfItemsToCreate = 500;
50+
51+
public static void main(String[] args) {
52+
SampleBulkHandleRetriesAsync p = new SampleBulkHandleRetriesAsync();
53+
54+
try {
55+
logger.info("Starting ASYNC main");
56+
p.getStartedDemo();
57+
logger.info("Demo complete, please hold while resources are released");
58+
} catch (Exception e) {
59+
e.printStackTrace();
60+
logger.error(String.format("Cosmos getStarted failed with %s", e));
61+
} finally {
62+
logger.info("Closing the client");
63+
p.shutdown();
64+
}
65+
}
66+
67+
public void close() {
68+
client.close();
69+
}
70+
71+
private void getStartedDemo() {
72+
73+
logger.info("Using Azure Cosmos DB endpoint: " + AccountSettings.HOST);
74+
75+
/* Create an async client
76+
force rate limiting by reducing the max retry attempts
77+
to simulate throttling when under heavy load*/
78+
client = new CosmosClientBuilder()
79+
.endpoint(AccountSettings.HOST)
80+
.credential(new DefaultAzureCredentialBuilder().build())
81+
.throttlingRetryOptions(
82+
new ThrottlingRetryOptions()
83+
.setMaxRetryAttemptsOnThrottledRequests(1)
84+
.setMaxRetryWaitTime(Duration.ofSeconds(1)))
85+
.contentResponseOnWriteEnabled(true)
86+
.consistencyLevel(ConsistencyLevel.SESSION).buildAsyncClient();
87+
88+
89+
createDatabaseIfNotExists();
90+
createContainerIfNotExists();
91+
92+
// Create 500 Family records
93+
ArrayList<Family> largeFamilies1 = new ArrayList<>();
94+
for (int i = 0; i < noOfItemsToCreate; i++) {
95+
Family family = new Family();
96+
family.setId(UUID.randomUUID().toString());
97+
family.setLastName("Family0");
98+
family.setRegistered(false);
99+
largeFamilies1.add(family);
100+
}
101+
logger.info("Ensure rate limiting with enough bulk upserts, retries handled by BulkWriter abstraction");
102+
largeBulkUpsertItemsWithBulkWriterAbstraction(largeFamilies1);
103+
}
104+
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+
119+
private void createContainerIfNotExists() {
120+
logger.info("Create container " + containerName + " if not exists.");
121+
122+
// Create container if not exists
123+
// <CreateContainerIfNotExists>
124+
125+
CosmosContainerProperties containerProperties = new CosmosContainerProperties(
126+
containerName, "/lastName");
127+
ThroughputProperties throughputProperties = ThroughputProperties.createManualThroughput(400);
128+
Mono<CosmosContainerResponse> containerIfNotExists = database
129+
.createContainerIfNotExists(containerProperties, throughputProperties);
130+
131+
// Create container with 400 RU/s
132+
CosmosContainerResponse cosmosContainerResponse = containerIfNotExists.block();
133+
assert(cosmosContainerResponse != null);
134+
assert(cosmosContainerResponse.getProperties() != null);
135+
container = database.getContainer(cosmosContainerResponse.getProperties().getId());
136+
// </CreateContainerIfNotExists>
137+
138+
//Modify existing container
139+
containerProperties = cosmosContainerResponse.getProperties();
140+
Mono<CosmosContainerResponse> propertiesReplace =
141+
container.replace(containerProperties, new CosmosContainerRequestOptions());
142+
propertiesReplace.flatMap(containerResponse -> {
143+
logger.info(
144+
"setupContainer(): Container {}} in {} has been updated with it's new properties.",
145+
container.getId(),
146+
database.getId());
147+
return Mono.empty();
148+
}).onErrorResume((exception) -> {
149+
logger.error(
150+
"setupContainer(): Unable to update properties for container {} in database {}. e: {}",
151+
container.getId(),
152+
database.getId(),
153+
exception.getLocalizedMessage(),
154+
exception);
155+
return Mono.empty();
156+
}).block();
157+
158+
}
159+
160+
private void largeBulkUpsertItemsWithBulkWriterAbstraction(Iterable<Family> families) {
161+
List<CosmosItemOperation> cosmosItemOperations = new ArrayList<>();
162+
for (Family family : families) {
163+
cosmosItemOperations.add(CosmosBulkOperations.getUpsertItemOperation(family, new PartitionKey(family.getLastName())));
164+
}
165+
BulkWriter bulkWriter = new BulkWriter(container);
166+
for (CosmosItemOperation operation : cosmosItemOperations) {
167+
bulkWriter.scheduleWrites(operation);
168+
}
169+
bulkWriter.execute().subscribe();
170+
//get count of items in container
171+
try {
172+
logger.info("Waiting for bulk operations to complete...");
173+
Thread.sleep(20000); // Wait for the bulk operations to complete
174+
} catch (InterruptedException e) {
175+
throw new RuntimeException(e);
176+
}
177+
logger.info("Number of items to create was: " + noOfItemsToCreate);
178+
logger.info("Total items created after bulk load: " + container.readAllItems(new PartitionKey("Family0"), Family.class).count().block());
179+
}
180+
181+
182+
private void shutdown() {
183+
try {
184+
// To allow for the sequence to complete after subscribe() calls
185+
Thread.sleep(5000);
186+
//Clean shutdown
187+
logger.info("Deleting Cosmos DB resources");
188+
logger.info("-Deleting container...");
189+
if (container != null) container.delete().subscribe();
190+
logger.info("-Deleting database...");
191+
if (database != null) database.delete().subscribe();
192+
logger.info("-Closing the client...");
193+
} catch (InterruptedException err) {
194+
err.printStackTrace();
195+
} catch (Exception err) {
196+
logger.error("Deleting Cosmos DB resources failed, will still attempt to close the client. See stack " + "trace below.");
197+
err.printStackTrace();
198+
}
199+
client.close();
200+
logger.info("Done.");
201+
}
202+
203+
204+
}

0 commit comments

Comments
 (0)