From a9637f1bbde53aa715cdce76f5e3dacac371672b Mon Sep 17 00:00:00 2001 From: Finn Zhang Date: Tue, 25 Aug 2026 21:47:57 +0000 Subject: [PATCH 1/7] feat: add sample codes and sample ITs for Cloud Spanner Queues. --- .../java/com/example/spanner/QueueSample.java | 188 ++++++++++++++++++ .../com/example/spanner/QueueSampleIT.java | 49 +++++ 2 files changed, 237 insertions(+) create mode 100644 java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java create mode 100644 java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java diff --git a/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java new file mode 100644 index 000000000000..55ece7ce9815 --- /dev/null +++ b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java @@ -0,0 +1,188 @@ +/* + * Copyright 2025 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.example.spanner; + +import com.google.cloud.ByteArray; +import com.google.cloud.Timestamp; +import com.google.cloud.spanner.Database; +import com.google.cloud.spanner.DatabaseAdminClient; +import com.google.cloud.spanner.DatabaseClient; +import com.google.cloud.spanner.DatabaseId; +import com.google.cloud.spanner.Key; +import com.google.cloud.spanner.Mutation; +import com.google.cloud.spanner.ResultSet; +import com.google.cloud.spanner.Spanner; +import com.google.cloud.spanner.SpannerOptions; +import com.google.cloud.spanner.Statement; +import com.google.cloud.spanner.Value; +import java.time.Duration; +import java.time.Instant; +import java.util.Arrays; +import java.util.Collections; +import java.util.concurrent.ExecutionException; + +public class QueueSample { + + static void createQueueDatabase(DatabaseAdminClient dbAdminClient, String instanceId, String databaseId) { + try { + System.out.println("Creating database with a queue..."); + Database database = + dbAdminClient + .createDatabase( + instanceId, + databaseId, + Collections.singletonList( + "CREATE Queue MyQueue (" + + " Id INT64 NOT NULL," + + " Payload BYTES(MAX) NOT NULL," + + ") PRIMARY KEY (Id), OPTIONS (receive_mode = 'PULL')")) + .get(); + System.out.println("Created database [" + database.getId() + "]"); + } catch (ExecutionException | InterruptedException e) { + System.err.println("Database creation failed: " + e.getMessage()); + } + } + + static void sendToQueueWithMutation(DatabaseClient dbClient) { + System.out.println("Sending a message to queue using Mutation API..."); + dbClient.write( + Collections.singletonList( + Mutation.newSendBuilder("MyQueue") + .setKey(Key.of(1L)) + .setPayload(Value.bytes(ByteArray.copyFrom("message1"))) + .build())); + System.out.println("Message sent."); + } + + static void sendToQueueWithSql(DatabaseClient dbClient) { + System.out.println("Sending a message to queue using SQL API..."); + dbClient.readWriteTransaction().run( + transaction -> { + transaction.executeUpdate( + Statement.of("INSERT INTO MyQueue (Id, Payload) VALUES (2, B'message2')")); + return null; + }); + System.out.println("Message sent."); + } + + static void sendToQueueWithMutationInFuture(DatabaseClient dbClient) { + System.out.println("Sending a message to queue using Mutation API in the future..."); + Instant futureTime = Instant.now().plus(Duration.ofMinutes(10)); + dbClient.write( + Collections.singletonList( + Mutation.newSendBuilder("MyQueue") + .setKey(Key.of(3L)) + .setPayload(Value.bytes(ByteArray.copyFrom("message3"))) + .setDeliveryTime(futureTime) + .build())); + System.out.println("Message scheduled for future delivery."); + } + + static void sendToQueueWithSqlInFuture(DatabaseClient dbClient) { + System.out.println("Sending a message to queue using SQL API in the future..."); + Instant futureTime = Instant.now().plus(Duration.ofMinutes(10)); + dbClient.readWriteTransaction().run( + transaction -> { + transaction.executeUpdate( + Statement.newBuilder("INSERT INTO MyQueue (Id, Payload, DeliverTime) VALUES (4, B'message4', @deliveryTime)") + .bind("deliveryTime").to(Value.timestamp(Timestamp.ofTimeSecondsAndNanos(futureTime.getEpochSecond(), futureTime.getNano()))) + .build()); + return null; + }); + System.out.println("Message scheduled for future delivery."); + } + + static void ackMessageWithMutation(DatabaseClient dbClient) { + System.out.println("Acknowledging a message using Mutation API..."); + dbClient.write( + Collections.singletonList( + Mutation.newAckBuilder("MyQueue") + .setKey(Key.of(1L)) + .build())); + System.out.println("Message acknowledged."); + } + + static void ackMessageWithSql(DatabaseClient dbClient) { + System.out.println("Acknowledging a message using SQL API with ASSERT_ROWS_MODIFIED..."); + dbClient.readWriteTransaction().run( + transaction -> { + transaction.executeUpdate( + Statement.of("DELETE FROM MyQueue WHERE Id = 2 ASSERT_ROWS_MODIFIED 1")); + return null; + }); + System.out.println("Message acknowledged."); + } + + static void deleteMessageWithSql(DatabaseClient dbClient) { + System.out.println("Deleting a message using SQL API..."); + dbClient.readWriteTransaction().run( + transaction -> { + transaction.executeUpdate( + Statement.of("DELETE FROM MyQueue WHERE Id = 3")); + return null; + }); + System.out.println("Message deleted."); + } + + static void sendAndReceiveWithSql(DatabaseClient dbClient) { + System.out.println("Sending and receiving a message using SQL API..."); + dbClient.readWriteTransaction().run( + transaction -> { + transaction.executeUpdate( + Statement.of("INSERT INTO MyQueue (Id, Payload) VALUES (5, B'message5')")); + return null; + }); + + System.out.println("Receiving message from queue (max_duration 1min)..."); + try (ResultSet resultSet = dbClient.singleUse().executeQuery( + Statement.of("SELECT * FROM RECEIVE_MyQueue(max_duration => '1m')"))) { + if (resultSet.next()) { + System.out.println("Received message ID: " + resultSet.getLong("Id")); + } else { + System.out.println("No messages received."); + } + } + } + + public static void main(String[] args) throws Exception { + if (args.length != 3) { + System.err.println("Usage: QueueSample "); + return; + } + String projectId = args[0]; + String instanceId = args[1]; + String databaseId = args[2]; + + try (Spanner spanner = SpannerOptions.newBuilder().setProjectId(projectId).build().getService()) { + DatabaseAdminClient dbAdminClient = spanner.getDatabaseAdminClient(); + createQueueDatabase(dbAdminClient, instanceId, databaseId); + + DatabaseClient dbClient = spanner.getDatabaseClient(DatabaseId.of(projectId, instanceId, databaseId)); + + sendToQueueWithMutation(dbClient); + sendToQueueWithSql(dbClient); + sendToQueueWithMutationInFuture(dbClient); + sendToQueueWithSqlInFuture(dbClient); + + ackMessageWithMutation(dbClient); + ackMessageWithSql(dbClient); + deleteMessageWithSql(dbClient); + + sendAndReceiveWithSql(dbClient); + } + } +} diff --git a/java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java b/java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java new file mode 100644 index 000000000000..f2f619376007 --- /dev/null +++ b/java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java @@ -0,0 +1,49 @@ +/* + * Copyright 2025 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.example.spanner; + +import static org.junit.Assert.assertTrue; + +import org.junit.Test; + +public class QueueSampleIT extends SampleTestBase { + + @Test + public void testQueueSample() throws Exception { + final String databaseId = idGenerator.generateDatabaseId(); + + final String out = + SampleRunner.runSample( + () -> + QueueSample.main( + new String[] { + projectId, instanceId, databaseId + })); + + assertTrue(out.contains("Creating database with a queue...")); + assertTrue(out.contains("Created database [")); + assertTrue(out.contains("Sending a message to queue using Mutation API...")); + assertTrue(out.contains("Sending a message to queue using SQL API...")); + assertTrue(out.contains("Sending a message to queue using Mutation API in the future...")); + assertTrue(out.contains("Sending a message to queue using SQL API in the future...")); + assertTrue(out.contains("Acknowledging a message using Mutation API...")); + assertTrue(out.contains("Acknowledging a message using SQL API with ASSERT_ROWS_MODIFIED...")); + assertTrue(out.contains("Deleting a message using SQL API...")); + assertTrue(out.contains("Sending and receiving a message using SQL API...")); + assertTrue(out.contains("Received message ID: 5")); + } +} From 521eba3b51cce2091899344553d95bca7455918b Mon Sep 17 00:00:00 2001 From: Finn Zhang Date: Thu, 27 Aug 2026 19:53:54 +0000 Subject: [PATCH 2/7] Fix copyright comments --- .../snippets/src/main/java/com/example/spanner/QueueSample.java | 2 +- .../src/test/java/com/example/spanner/QueueSampleIT.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java index 55ece7ce9815..7ebcb04f2d0e 100644 --- a/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java +++ b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java @@ -1,5 +1,5 @@ /* - * Copyright 2025 Google LLC + * Copyright 2026 Google LLC * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. diff --git a/java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java b/java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java index f2f619376007..38136723503e 100644 --- a/java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java +++ b/java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java @@ -1,5 +1,5 @@ /* - * Copyright 2025 Google LLC + * Copyright 2026 Google LLC * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. From 0b2b6210b0334bca7d650bbf596a4456528d3639 Mon Sep 17 00:00:00 2001 From: Finn Zhang Date: Mon, 31 Aug 2026 23:28:11 +0000 Subject: [PATCH 3/7] Fix gemini comment --- .../java/com/example/spanner/QueueSample.java | 33 +++++++++---------- 1 file changed, 15 insertions(+), 18 deletions(-) diff --git a/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java index 7ebcb04f2d0e..2a347efd4079 100644 --- a/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java +++ b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java @@ -37,24 +37,21 @@ public class QueueSample { - static void createQueueDatabase(DatabaseAdminClient dbAdminClient, String instanceId, String databaseId) { - try { - System.out.println("Creating database with a queue..."); - Database database = - dbAdminClient - .createDatabase( - instanceId, - databaseId, - Collections.singletonList( - "CREATE Queue MyQueue (" - + " Id INT64 NOT NULL," - + " Payload BYTES(MAX) NOT NULL," - + ") PRIMARY KEY (Id), OPTIONS (receive_mode = 'PULL')")) - .get(); - System.out.println("Created database [" + database.getId() + "]"); - } catch (ExecutionException | InterruptedException e) { - System.err.println("Database creation failed: " + e.getMessage()); - } + static void createQueueDatabase(DatabaseAdminClient dbAdminClient, String instanceId, String databaseId) + throws ExecutionException, InterruptedException { + System.out.println("Creating database with a queue..."); + Database database = + dbAdminClient + .createDatabase( + instanceId, + databaseId, + Collections.singletonList( + "CREATE Queue MyQueue (" + + " Id INT64 NOT NULL," + + " Payload BYTES(MAX) NOT NULL," + + ") PRIMARY KEY (Id), OPTIONS (receive_mode = 'PULL')")) + .get(); + System.out.println("Created database [" + database.getId() + "]"); } static void sendToQueueWithMutation(DatabaseClient dbClient) { From 0ccfc51bdbc9902385c37de6d267b0f59af14db3 Mon Sep 17 00:00:00 2001 From: Finn Zhang Date: Tue, 1 Sep 2026 20:30:05 +0000 Subject: [PATCH 4/7] Skip Emulator in IT --- .../src/test/java/com/example/spanner/QueueSampleIT.java | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java b/java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java index 38136723503e..04e9d07ced01 100644 --- a/java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java +++ b/java-spanner/samples/snippets/src/test/java/com/example/spanner/QueueSampleIT.java @@ -17,11 +17,20 @@ package com.example.spanner; import static org.junit.Assert.assertTrue; +import static org.junit.Assume.assumeTrue; +import org.junit.BeforeClass; import org.junit.Test; public class QueueSampleIT extends SampleTestBase { + @BeforeClass + public static void setUpTestSuite() { + assumeTrue( + "Queue is not supported by the Spanner Emulator", + System.getenv("SPANNER_EMULATOR_HOST") == null); + } + @Test public void testQueueSample() throws Exception { final String databaseId = idGenerator.generateDatabaseId(); From a004f086c501110d854cbb13a4c316d8c00fb942 Mon Sep 17 00:00:00 2001 From: Finn Zhang Date: Fri, 4 Sep 2026 05:42:19 +0000 Subject: [PATCH 5/7] add sample tags --- .../src/main/java/com/example/spanner/QueueSample.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java index 2a347efd4079..4cae6babc63b 100644 --- a/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java +++ b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java @@ -16,6 +16,8 @@ package com.example.spanner; +// [START spanner_queue_sample] + import com.google.cloud.ByteArray; import com.google.cloud.Timestamp; import com.google.cloud.spanner.Database; @@ -183,3 +185,4 @@ public static void main(String[] args) throws Exception { } } } +// [END spanner_queue_sample] From 9aae8adadf86f0bc1793729871a933b10c79190b Mon Sep 17 00:00:00 2001 From: Finn Zhang Date: Fri, 4 Sep 2026 06:41:54 +0000 Subject: [PATCH 6/7] Update sample tags to match for different languages --- .../java/com/example/spanner/QueueSample.java | 30 ++++++++++++++----- 1 file changed, 22 insertions(+), 8 deletions(-) diff --git a/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java index 4cae6babc63b..90a3fc51377a 100644 --- a/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java +++ b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java @@ -1,5 +1,5 @@ /* - * Copyright 2026 Google LLC + * Copyright 2025 Google LLC * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -16,8 +16,6 @@ package com.example.spanner; -// [START spanner_queue_sample] - import com.google.cloud.ByteArray; import com.google.cloud.Timestamp; import com.google.cloud.spanner.Database; @@ -39,8 +37,8 @@ public class QueueSample { - static void createQueueDatabase(DatabaseAdminClient dbAdminClient, String instanceId, String databaseId) - throws ExecutionException, InterruptedException { +// [START spanner_create_database_with_queue] + static void createQueueDatabase(DatabaseAdminClient dbAdminClient, String instanceId, String databaseId) throws ExecutionException, InterruptedException { System.out.println("Creating database with a queue..."); Database database = dbAdminClient @@ -55,7 +53,9 @@ static void createQueueDatabase(DatabaseAdminClient dbAdminClient, String instan .get(); System.out.println("Created database [" + database.getId() + "]"); } +// [END spanner_create_database_with_queue] +// [START spanner_send_to_queue_with_mutation_api] static void sendToQueueWithMutation(DatabaseClient dbClient) { System.out.println("Sending a message to queue using Mutation API..."); dbClient.write( @@ -66,7 +66,9 @@ static void sendToQueueWithMutation(DatabaseClient dbClient) { .build())); System.out.println("Message sent."); } +// [END spanner_send_to_queue_with_mutation_api] +// [START spanner_send_to_queue_with_sql_api] static void sendToQueueWithSql(DatabaseClient dbClient) { System.out.println("Sending a message to queue using SQL API..."); dbClient.readWriteTransaction().run( @@ -77,7 +79,9 @@ static void sendToQueueWithSql(DatabaseClient dbClient) { }); System.out.println("Message sent."); } +// [END spanner_send_to_queue_with_sql_api] +// [START spanner_send_to_queue_with_mutation_api_in_future] static void sendToQueueWithMutationInFuture(DatabaseClient dbClient) { System.out.println("Sending a message to queue using Mutation API in the future..."); Instant futureTime = Instant.now().plus(Duration.ofMinutes(10)); @@ -90,21 +94,25 @@ static void sendToQueueWithMutationInFuture(DatabaseClient dbClient) { .build())); System.out.println("Message scheduled for future delivery."); } +// [END spanner_send_to_queue_with_mutation_api_in_future] +// [START spanner_send_to_queue_with_sql_api_in_future] static void sendToQueueWithSqlInFuture(DatabaseClient dbClient) { System.out.println("Sending a message to queue using SQL API in the future..."); Instant futureTime = Instant.now().plus(Duration.ofMinutes(10)); dbClient.readWriteTransaction().run( transaction -> { transaction.executeUpdate( - Statement.newBuilder("INSERT INTO MyQueue (Id, Payload, DeliverTime) VALUES (4, B'message4', @deliveryTime)") + Statement.newBuilder("INSERT INTO MyQueue (Id, Payload, enqueued_time) VALUES (4, B'message4', @deliveryTime)") .bind("deliveryTime").to(Value.timestamp(Timestamp.ofTimeSecondsAndNanos(futureTime.getEpochSecond(), futureTime.getNano()))) .build()); return null; }); System.out.println("Message scheduled for future delivery."); } +// [END spanner_send_to_queue_with_sql_api_in_future] +// [START spanner_ack_queue_message_with_mutation_api] static void ackMessageWithMutation(DatabaseClient dbClient) { System.out.println("Acknowledging a message using Mutation API..."); dbClient.write( @@ -114,7 +122,9 @@ static void ackMessageWithMutation(DatabaseClient dbClient) { .build())); System.out.println("Message acknowledged."); } +// [END spanner_ack_queue_message_with_mutation_api] +// [START spanner_ack_queue_message_with_sql_api] static void ackMessageWithSql(DatabaseClient dbClient) { System.out.println("Acknowledging a message using SQL API with ASSERT_ROWS_MODIFIED..."); dbClient.readWriteTransaction().run( @@ -125,7 +135,9 @@ static void ackMessageWithSql(DatabaseClient dbClient) { }); System.out.println("Message acknowledged."); } +// [END spanner_ack_queue_message_with_sql_api] +// [START spanner_delete_queue_message_with_sql_api] static void deleteMessageWithSql(DatabaseClient dbClient) { System.out.println("Deleting a message using SQL API..."); dbClient.readWriteTransaction().run( @@ -136,7 +148,9 @@ static void deleteMessageWithSql(DatabaseClient dbClient) { }); System.out.println("Message deleted."); } +// [END spanner_delete_queue_message_with_sql_api] +// [START spanner_send_and_receive_queue_message_with_sql_api] static void sendAndReceiveWithSql(DatabaseClient dbClient) { System.out.println("Sending and receiving a message using SQL API..."); dbClient.readWriteTransaction().run( @@ -148,7 +162,7 @@ static void sendAndReceiveWithSql(DatabaseClient dbClient) { System.out.println("Receiving message from queue (max_duration 1min)..."); try (ResultSet resultSet = dbClient.singleUse().executeQuery( - Statement.of("SELECT * FROM RECEIVE_MyQueue(max_duration => '1m')"))) { + Statement.of("SELECT * FROM READ_MyQueue(max_duration => '1m')"))) { if (resultSet.next()) { System.out.println("Received message ID: " + resultSet.getLong("Id")); } else { @@ -156,6 +170,7 @@ static void sendAndReceiveWithSql(DatabaseClient dbClient) { } } } +// [END spanner_send_and_receive_queue_message_with_sql_api] public static void main(String[] args) throws Exception { if (args.length != 3) { @@ -185,4 +200,3 @@ public static void main(String[] args) throws Exception { } } } -// [END spanner_queue_sample] From 3c6c2bd7ee34ec44eb7ec84ed56620eac471ad27 Mon Sep 17 00:00:00 2001 From: Finn Zhang Date: Fri, 4 Sep 2026 06:43:53 +0000 Subject: [PATCH 7/7] fix header --- .../snippets/src/main/java/com/example/spanner/QueueSample.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java index 90a3fc51377a..cd95a3fd8779 100644 --- a/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java +++ b/java-spanner/samples/snippets/src/main/java/com/example/spanner/QueueSample.java @@ -1,5 +1,5 @@ /* - * Copyright 2025 Google LLC + * Copyright 2026 Google LLC * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License.