Skip to content

Commit 4d5bd06

Browse files
authored
Merge pull request #473 from weaviate/v6-backup
v6: Backups
2 parents 60c03e9 + 6dbaa54 commit 4d5bd06

26 files changed

Lines changed: 1431 additions & 8 deletions

README.md

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -733,6 +733,64 @@ client.collections.update("Songs_Alias", "PopSongs");
733733
client.collections.delete("Songs_Alias");
734734
```
735735
736+
### Managing collection backups
737+
738+
> [!CAUTION]
739+
> Weaviate does not support concurrent backups. Await one backup's completion before starting another one.
740+
741+
```java
742+
// Start a backup:
743+
Backup backup = client.backup.create(
744+
"backup_1", "filesystem",
745+
bak -> bak
746+
.includeCollections("Songs", "Artists")
747+
.compressionLevel(CompressionLevel.BEST_COMPRESSION)
748+
.cpuPercentage(30)
749+
);
750+
751+
// By default, the client does not monitor the backup status.
752+
// The above method returns as soon as the server acknowledges
753+
// the request and starts to process it.
754+
//
755+
// Now you can poll backup status to know when it is succeedes (or fails).
756+
757+
Backup status = client.backup.getCreateStatus(backup.id(), backup.backend());
758+
if (status.status() == BackupStatus.SUCCESSFUL) {
759+
System.out.println("Yay!");
760+
System.exit(0);
761+
}
762+
763+
// Backups may take a write to complete. To block the current thread until
764+
// the execution completes, call Backup::waitForCompletion(WeaviateClient).
765+
//
766+
// Notice that, while we use `backup` object we can also call it on the `status`,
767+
// as both will have sufficient information to identify the backup operation.
768+
769+
try {
770+
Backup completed = backup.waitForCompletion(client);
771+
assert completed.errors() == null : "completed with errors";
772+
} catch (TimeoutException e) {
773+
System.out.exit(1);
774+
}
775+
776+
// List exists backups:
777+
List<Backup> allBackups = client.backup.list();
778+
779+
// Restore from the first backup:
780+
var first = allBackups.getFirst();
781+
client.backup.restore(first.id(), first.backend());
782+
783+
// Similarly, wait until the restore is complete using Backup::waitForCompletion.
784+
// It is possible to set a custom timeout and polling interval using a familiar Tucked Builder pattern:
785+
786+
var restoring = client.backup.getRestoreStatus(first.id(), first.backend());
787+
var restored = restoring.waitForCompletion(client, wait -> wait
788+
.timeout(Duration.ofMinutes(30))
789+
.interval(Duration.ofMinutes(5)));
790+
791+
assert restored.errors() == null : "restored with errors";
792+
```
793+
736794
### RBAC
737795
738796
#### Roles

src/it/java/io/weaviate/containers/Weaviate.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -133,6 +133,11 @@ public Builder withOffloadS3(String accessKey, String secretKey) {
133133
return this;
134134
}
135135

136+
public Builder withFilesystemBackup(String fsPath) {
137+
addModules("backup-filesystem");
138+
environment.put("BACKUP_FILESYSTEM_PATH", fsPath);
139+
return this;
140+
}
136141
public Builder withAdminUsers(String... admins) {
137142
adminUsers.addAll(Arrays.asList(admins));
138143
return this;

src/it/java/io/weaviate/integration/AliasITest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@
1212
import io.weaviate.containers.Container;
1313

1414
public class AliasITest extends ConcurrentTest {
15-
private static WeaviateClient client = Container.WEAVIATE.getClient();
15+
private static final WeaviateClient client = Container.WEAVIATE.getClient();
1616

1717
@Test
1818
public void test_aliasLifecycle() throws IOException {
Lines changed: 251 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,251 @@
1+
package io.weaviate.integration;
2+
3+
import java.io.IOException;
4+
import java.util.List;
5+
import java.util.Map;
6+
import java.util.UUID;
7+
import java.util.concurrent.CompletableFuture;
8+
import java.util.concurrent.CompletionException;
9+
import java.util.concurrent.ExecutionException;
10+
import java.util.concurrent.TimeoutException;
11+
import java.util.stream.IntStream;
12+
13+
import org.assertj.core.api.Assertions;
14+
import org.assertj.core.api.InstanceOfAssertFactories;
15+
import org.junit.Test;
16+
17+
import io.weaviate.ConcurrentTest;
18+
import io.weaviate.client6.v1.api.WeaviateClient;
19+
import io.weaviate.client6.v1.api.backup.Backup;
20+
import io.weaviate.client6.v1.api.backup.BackupStatus;
21+
import io.weaviate.client6.v1.api.backup.CompressionLevel;
22+
import io.weaviate.client6.v1.api.collections.Property;
23+
import io.weaviate.containers.Weaviate;
24+
25+
public class BackupITest extends ConcurrentTest {
26+
private static final WeaviateClient client = Weaviate.custom()
27+
.withFilesystemBackup("/tmp/backups").build()
28+
.getClient();
29+
30+
@Test
31+
public void test_lifecycle() throws IOException, TimeoutException {
32+
// Arrange
33+
String nsA = ns("A"), nsB = ns("B"), nsC = ns("C"), nsBig = ns("Big");
34+
String backup_1 = ns("backup_1").toLowerCase();
35+
String backend = "filesystem";
36+
37+
// Start writing data in the background so it's ready
38+
// by the time we get to backup #3.
39+
var spam = spamData(nsBig);
40+
41+
client.collections.create(nsA);
42+
client.collections.create(nsB);
43+
client.collections.create(nsC);
44+
45+
// Insert some data to check restore later
46+
var collectionA = client.collections.use(nsA);
47+
collectionA.data.insert(Map.of());
48+
49+
// Act: start backup
50+
var started = client.backup.create(backup_1, backend,
51+
backup -> backup
52+
.includeCollections(nsA, nsB)
53+
.compressionLevel(CompressionLevel.BEST_SPEED));
54+
55+
// Assert
56+
Assertions.assertThat(started)
57+
.as("created backup operation")
58+
.returns(backup_1, Backup::id)
59+
.returns(backend, Backup::backend)
60+
.returns(BackupStatus.STARTED, Backup::status)
61+
.returns(null, Backup::error)
62+
.extracting(Backup::includesCollections, InstanceOfAssertFactories.list(String.class))
63+
.containsOnly(nsA, nsB);
64+
65+
// Act: await backup competion
66+
var completed = started.waitForCompletion(client);
67+
68+
// Assert
69+
Assertions.assertThat(completed)
70+
.as("await backup completion")
71+
.returns(backup_1, Backup::id)
72+
.returns(backend, Backup::backend)
73+
.returns(BackupStatus.SUCCESS, Backup::status)
74+
.returns(null, Backup::error);
75+
76+
// Act: create another backup
77+
String backup_2 = ns("backup_2").toLowerCase();
78+
client.backup.create(backup_2, backend).waitForCompletion(client);
79+
80+
// Assert: check the second backup is created successfully
81+
var status_2 = client.backup.getCreateStatus(backup_2, backend);
82+
Assertions.assertThat(status_2).as("backup #2").get()
83+
.returns(BackupStatus.SUCCESS, Backup::status);
84+
85+
// Act: create and cancel
86+
// Try to throttle this backup by creating a lot of objects,
87+
// limiting CPU resources and requiring high compression ratio.
88+
// This is to avoid flaky tests and make sure we can cancel
89+
// the backup before it completes successfully.
90+
String backup_3 = ns("backup_3").toLowerCase();
91+
spam.join();
92+
var cancelMe = client.backup.create(backup_3, backend,
93+
backup -> backup
94+
.includeCollections(nsA, nsB, nsC, nsBig)
95+
.cpuPercentage(1)
96+
.compressionLevel(CompressionLevel.BEST_COMPRESSION));
97+
cancelMe.cancel(client);
98+
cancelMe.waitForStatus(client, BackupStatus.CANCELED, wait -> wait.interval(500));
99+
100+
// Assert: check the backup is cancelled
101+
var status_3 = client.backup.getCreateStatus(backup_3, backend);
102+
Assertions.assertThat(status_3).as("backup #3").get()
103+
.returns(BackupStatus.CANCELED, Backup::status);
104+
105+
// Assert: all 3 backups are present
106+
var all = client.backup.list(backend);
107+
Assertions.assertThat(all).as("all backups")
108+
.extracting(Backup::id)
109+
.contains(backup_1, backup_2, backup_3);
110+
111+
// Act: delete data and restore backup #1
112+
client.collections.delete(nsA);
113+
client.backup.restore(backup_1, backend, restore -> restore.includeCollections(nsA));
114+
115+
// Assert: object inserted in the beginning of the test is present
116+
var restore_1 = client.backup.getRestoreStatus(backup_1, backend)
117+
.orElseThrow().waitForCompletion(client);
118+
Assertions.assertThat(restore_1).as("restore backup #1")
119+
.returns(BackupStatus.SUCCESS, Backup::status);
120+
Assertions.assertThat(collectionA.size()).as("after restore backup #1").isEqualTo(1);
121+
}
122+
123+
@Test
124+
public void test_lifecycle_async() throws ExecutionException, InterruptedException, Exception {
125+
// Arrange
126+
String nsA = ns("A"), nsB = ns("B"), nsC = ns("C"), nsBig = ns("Big");
127+
String backup_1 = ns("backup_1").toLowerCase();
128+
String backend = "filesystem";
129+
130+
try (final var async = client.async()) {
131+
// Start writing data in the background so it's ready
132+
// by the time we get to backup #3.
133+
var spam = spamData(nsBig);
134+
135+
CompletableFuture.allOf(
136+
async.collections.create(nsA),
137+
async.collections.create(nsB),
138+
async.collections.create(nsC))
139+
.join();
140+
141+
// Insert some data to check restore later
142+
var collectionA = async.collections.use(nsA);
143+
collectionA.data.insert(Map.of()).join();
144+
145+
// Act: start backup
146+
var started = async.backup.create(backup_1, backend,
147+
backup -> backup
148+
.includeCollections(nsA, nsB)
149+
.compressionLevel(CompressionLevel.BEST_SPEED))
150+
.join();
151+
152+
// Assert
153+
Assertions.assertThat(started)
154+
.as("created backup operation")
155+
.returns(backup_1, Backup::id)
156+
.returns(backend, Backup::backend)
157+
.returns(BackupStatus.STARTED, Backup::status)
158+
.returns(null, Backup::error)
159+
.extracting(Backup::includesCollections, InstanceOfAssertFactories.list(String.class))
160+
.containsOnly(nsA, nsB);
161+
162+
// Act: await backup competion
163+
var completed = started.waitForCompletion(async).join();
164+
165+
// Assert
166+
Assertions.assertThat(completed)
167+
.as("await backup completion")
168+
.returns(backup_1, Backup::id)
169+
.returns(backend, Backup::backend)
170+
.returns(BackupStatus.SUCCESS, Backup::status)
171+
.returns(null, Backup::error);
172+
173+
// Act: create another backup
174+
String backup_2 = ns("backup_2").toLowerCase();
175+
async.backup.create(backup_2, backend)
176+
.thenCompose(bak -> bak.waitForCompletion(async))
177+
.join();
178+
179+
// Assert: check the second backup is created successfully
180+
var status_2 = async.backup.getCreateStatus(backup_2, backend).join();
181+
Assertions.assertThat(status_2).as("backup #2").get()
182+
.returns(BackupStatus.SUCCESS, Backup::status);
183+
184+
// Act: create and cancel
185+
// Try to throttle this backup by creating a lot of objects,
186+
// limiting CPU resources and requiring high compression ratio.
187+
// This is to avoid flaky tests and make sure we can cancel
188+
// the backup before it completes successfully.
189+
String backup_3 = ns("backup_3").toLowerCase();
190+
spam.join();
191+
async.backup.create(backup_3, backend,
192+
backup -> backup
193+
.includeCollections(nsA, nsB, nsC, nsBig)
194+
.cpuPercentage(1)
195+
.compressionLevel(CompressionLevel.BEST_COMPRESSION))
196+
.thenCompose(cancelMe -> cancelMe.cancel(async)
197+
.thenCompose(__ -> cancelMe.waitForStatus(async, BackupStatus.CANCELED,
198+
wait -> wait.interval(500))))
199+
.join();
200+
201+
// Assert: check the backup is cancelled
202+
var status_3 = async.backup.getCreateStatus(backup_3, backend).join();
203+
Assertions.assertThat(status_3).as("backup #3").get()
204+
.returns(BackupStatus.CANCELED, Backup::status);
205+
206+
// Assert: all 3 backups are present
207+
var all = async.backup.list(backend).join();
208+
Assertions.assertThat(all).as("all backups")
209+
.extracting(Backup::id)
210+
.contains(backup_1, backup_2, backup_3);
211+
212+
// Act: delete data and restore backup #1
213+
async.collections.delete(nsA).join();
214+
async.backup.restore(backup_1, backend, restore -> restore.includeCollections(nsA)).join();
215+
216+
// Assert: object inserted in the beginning of the test is present
217+
var restore_1 = async.backup.getRestoreStatus(backup_1, backend)
218+
.thenCompose(bak -> bak.orElseThrow().waitForCompletion(async))
219+
.join();
220+
Assertions.assertThat(restore_1).as("restore backup #1")
221+
.returns(BackupStatus.SUCCESS, Backup::status);
222+
Assertions.assertThat(collectionA.size().join()).as("after restore backup #1").isEqualTo(1);
223+
}
224+
}
225+
226+
@Test(expected = IllegalStateException.class)
227+
public void test_waitForCompletion_unknown() throws IOException, TimeoutException {
228+
var backup = new Backup("#1", "/tmp/bak/#1", "filesystem", List.of("Things"), BackupStatus.STARTED, null,
229+
null);
230+
backup.waitForCompletion(client);
231+
}
232+
233+
/** Write 10_000 entries with a UUID[10] property. */
234+
private CompletableFuture<Void> spamData(String collectionName) {
235+
return CompletableFuture.supplyAsync(() -> {
236+
try {
237+
client.collections.create(collectionName,
238+
c -> c.properties(Property.uuidArray("uuids")));
239+
240+
var spam = client.collections.use(collectionName);
241+
for (int i = 0; i < 10_000; i++) {
242+
var uuids = IntStream.range(0, 10).mapToObj(j -> UUID.randomUUID()).toArray();
243+
spam.data.insert(Map.of("uuids", uuids));
244+
}
245+
} catch (IOException e) {
246+
throw new CompletionException(e);
247+
}
248+
return null;
249+
});
250+
}
251+
}

src/it/java/io/weaviate/integration/DataITest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,7 @@
2929
import io.weaviate.containers.Container;
3030

3131
public class DataITest extends ConcurrentTest {
32-
private static WeaviateClient client = Container.WEAVIATE.getClient();
32+
private static final WeaviateClient client = Container.WEAVIATE.getClient();
3333
private static final String COLLECTION = unique("Artists");
3434
private static final String VECTOR_INDEX = "bring_your_own";
3535

src/it/java/io/weaviate/integration/ORMITest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@
2323
import io.weaviate.containers.Container;
2424

2525
public class ORMITest extends ConcurrentTest {
26-
private static WeaviateClient client = Container.WEAVIATE.getClient();
26+
private static final WeaviateClient client = Container.WEAVIATE.getClient();
2727

2828
@Collection("ORMITestThings")
2929
static record Thing(

src/it/java/io/weaviate/integration/PaginationITest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@
2525
import io.weaviate.containers.Container;
2626

2727
public class PaginationITest extends ConcurrentTest {
28-
private static WeaviateClient client = Container.WEAVIATE.getClient();
28+
private static final WeaviateClient client = Container.WEAVIATE.getClient();
2929

3030
@Test
3131
public void testIterateAll() throws IOException {

src/it/java/io/weaviate/integration/TenantsITest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@ public class TenantsITest extends ConcurrentTest {
1818
.build(),
1919
Container.MINIO);
2020

21-
private static WeaviateClient client = compose.getClient();
21+
private static final WeaviateClient client = compose.getClient();
2222

2323
@Test
2424
public void test_tenantLifecycle() throws Exception {

src/main/java/io/weaviate/client6/v1/api/WeaviateClient.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44
import java.util.function.Function;
55

66
import io.weaviate.client6.v1.api.alias.WeaviateAliasClient;
7+
import io.weaviate.client6.v1.api.backup.WeaviateBackupClient;
78
import io.weaviate.client6.v1.api.collections.WeaviateCollectionsClient;
89
import io.weaviate.client6.v1.api.rbac.groups.WeaviateGroupsClient;
910
import io.weaviate.client6.v1.api.rbac.roles.WeaviateRolesClient;
@@ -34,6 +35,9 @@ public class WeaviateClient implements AutoCloseable {
3435
/** Client for {@code /aliases} endpoints for managing collection aliases. */
3536
public final WeaviateAliasClient alias;
3637

38+
/** Client for {@code /backups} endpoints for managing backups. */
39+
public final WeaviateBackupClient backup;
40+
3741
/**
3842
* Client for {@code /authz/roles} endpoints for managing RBAC roles.
3943
*/
@@ -99,6 +103,7 @@ public WeaviateClient(Config config) {
99103
this.restTransport = _restTransport;
100104
this.grpcTransport = new DefaultGrpcTransport(grpcOpt);
101105
this.alias = new WeaviateAliasClient(restTransport);
106+
this.backup = new WeaviateBackupClient(restTransport);
102107
this.collections = new WeaviateCollectionsClient(restTransport, grpcTransport);
103108
this.roles = new WeaviateRolesClient(restTransport);
104109
this.groups = new WeaviateGroupsClient(restTransport);

0 commit comments

Comments
 (0)