|
20 | 20 | from typing import List, Optional |
21 | 21 |
|
22 | 22 | from google.api_core.exceptions import InternalServerError |
| 23 | +from google.cloud._helpers import _datetime_to_pb_timestamp |
23 | 24 |
|
24 | 25 | from google.cloud.aio._cross_sync import CrossSync |
25 | 26 | from google.cloud.spanner_v1._async._helpers import _retry, _retry_on_aborted_exception |
26 | 27 | from google.cloud.spanner_v1._helpers import ( |
27 | 28 | AtomicCounter, |
28 | 29 | _check_rst_stream_error, |
| 30 | + _make_list_value_pb, |
29 | 31 | _make_list_value_pbs, |
| 32 | + _make_value_pb, |
30 | 33 | _merge_client_context, |
31 | 34 | _merge_request_options, |
32 | 35 | _merge_Transaction_Options, |
@@ -165,6 +168,49 @@ def delete(self, table, keyset): |
165 | 168 | # TODO: Decide if we should add a span event per mutation: |
166 | 169 | # https://github.com/googleapis/python-spanner/issues/1269 |
167 | 170 |
|
| 171 | + def send(self, queue, key, payload=None, deliver_time=None): |
| 172 | + """Send a message to a Cloud Spanner queue. |
| 173 | +
|
| 174 | + :type queue: str |
| 175 | + :param queue: Name of the queue to which the message will be sent. |
| 176 | +
|
| 177 | + :type key: list |
| 178 | + :param key: The primary key of the message to be sent. |
| 179 | +
|
| 180 | + :type payload: object |
| 181 | + :param payload: (Optional) The payload of the message. |
| 182 | +
|
| 183 | + :type deliver_time: :class:`datetime.datetime` |
| 184 | + :param deliver_time: (Optional) The time at which Spanner will begin attempting to deliver the message. |
| 185 | + """ |
| 186 | + send_kwargs = {"queue": queue, "key": _make_list_value_pb(key)} |
| 187 | + if payload is not None: |
| 188 | + send_kwargs["payload"] = _make_value_pb(payload) |
| 189 | + if deliver_time is not None: |
| 190 | + send_kwargs["deliver_time"] = _datetime_to_pb_timestamp(deliver_time) |
| 191 | + |
| 192 | + send = Mutation.Send(**send_kwargs) |
| 193 | + self._mutations.append(Mutation(send=send)) |
| 194 | + |
| 195 | + def ack(self, queue, key, ignore_not_found=None): |
| 196 | + """Acknowledge a message in a Cloud Spanner queue. |
| 197 | +
|
| 198 | + :type queue: str |
| 199 | + :param queue: Name of the queue where the message to be acked is stored. |
| 200 | +
|
| 201 | + :type key: list |
| 202 | + :param key: The primary key of the message to be acked. |
| 203 | +
|
| 204 | + :type ignore_not_found: bool |
| 205 | + :param ignore_not_found: (Optional) Whether to ignore if the message does not exist. |
| 206 | + """ |
| 207 | + ack_kwargs = {"queue": queue, "key": _make_list_value_pb(key)} |
| 208 | + if ignore_not_found is not None: |
| 209 | + ack_kwargs["ignore_not_found"] = ignore_not_found |
| 210 | + |
| 211 | + ack = Mutation.Ack(**ack_kwargs) |
| 212 | + self._mutations.append(Mutation(ack=ack)) |
| 213 | + |
168 | 214 |
|
169 | 215 | class Batch(_BatchBase): |
170 | 216 | """Accumulate mutations for transmission during :meth:`commit`.""" |
|
0 commit comments