Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 21 additions & 8 deletions google/cloud/spanner_v1/batch.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
_metadata_with_prefix,
_metadata_with_leader_aware_routing,
_merge_Transaction_Options,
AtomicCounter,
)
from google.cloud.spanner_v1._opentelemetry_tracing import trace_call
from google.cloud.spanner_v1 import RequestOptions
Expand Down Expand Up @@ -248,7 +249,7 @@ def commit(
trace_attributes,
observability_options=observability_options,
metadata=metadata,
), MetricsCapture():
) as span, MetricsCapture():

def wrapped_method(*args, **kwargs):
method = functools.partial(
Expand All @@ -261,6 +262,7 @@ def wrapped_method(*args, **kwargs):
getattr(database, "_next_nth_request", 0),
1,
metadata,
span,
),
)
return method(*args, **kwargs)
Expand Down Expand Up @@ -384,14 +386,25 @@ def batch_write(self, request_options=None, exclude_txn_from_change_streams=Fals
trace_attributes,
observability_options=observability_options,
metadata=metadata,
), MetricsCapture():
method = functools.partial(
api.batch_write,
request=request,
metadata=metadata,
)
) as span, MetricsCapture():
attempt = AtomicCounter(0)
nth_request = getattr(database, "_next_nth_request", 0)

def wrapped_method(*args, **kwargs):
method = functools.partial(
api.batch_write,
request=request,
metadata=database.metadata_with_request_id(
nth_request,
attempt.increment(),
metadata,
span,
),
)
return method(*args, **kwargs)

response = _retry(
method,
wrapped_method,
allowed_exceptions={
InternalServerError: _check_rst_stream_error,
},
Expand Down
9 changes: 8 additions & 1 deletion google/cloud/spanner_v1/database.py
Original file line number Diff line number Diff line change
Expand Up @@ -462,13 +462,19 @@ def spanner_api(self):

return self._spanner_api

def metadata_with_request_id(self, nth_request, nth_attempt, prior_metadata=[]):
def metadata_with_request_id(
self, nth_request, nth_attempt, prior_metadata=[], span=None
):
if span is None:
span = get_current_span()

return _metadata_with_request_id(
self._nth_client_id,
self._channel_id,
nth_request,
nth_attempt,
prior_metadata,
span,
)

def __eq__(self, other):
Expand Down Expand Up @@ -762,6 +768,7 @@ def execute_pdml():
self._next_nth_request,
1,
metadata,
span,
),
)

Expand Down
10 changes: 8 additions & 2 deletions google/cloud/spanner_v1/pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -257,7 +257,10 @@ def bind(self, database):
resp = api.batch_create_sessions(
request=request,
metadata=database.metadata_with_request_id(
database._next_nth_request, 1, metadata
database._next_nth_request,
1,
metadata,
span,
),
)

Expand Down Expand Up @@ -564,7 +567,10 @@ def bind(self, database):
resp = api.batch_create_sessions(
request=request,
metadata=database.metadata_with_request_id(
database._next_nth_request, 1, metadata
database._next_nth_request,
1,
metadata,
span,
),
)

Expand Down
24 changes: 23 additions & 1 deletion google/cloud/spanner_v1/request_id_header.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,10 +33,32 @@ def generate_rand_uint64():


REQ_RAND_PROCESS_ID = generate_rand_uint64()
X_GOOG_SPANNER_REQUEST_ID_SPAN_ATTR = "x_goog_spanner_request_id"


def with_request_id(client_id, channel_id, nth_request, attempt, other_metadata=[]):
def with_request_id(
client_id, channel_id, nth_request, attempt, other_metadata=[], span=None
):
req_id = f"{REQ_ID_VERSION}.{REQ_RAND_PROCESS_ID}.{client_id}.{channel_id}.{nth_request}.{attempt}"
all_metadata = (other_metadata or []).copy()
all_metadata.append((REQ_ID_HEADER_KEY, req_id))

if span is not None:
span.set_attribute(X_GOOG_SPANNER_REQUEST_ID_SPAN_ATTR, req_id)

return all_metadata


def parse_request_id(request_id_str):
splits = request_id_str.split(".")
version, rand_process_id, client_id, channel_id, nth_request, nth_attempt = list(
map(lambda v: int(v), splits)
)
return (
version,
rand_process_id,
client_id,
channel_id,
nth_request,
nth_attempt,
)
19 changes: 14 additions & 5 deletions google/cloud/spanner_v1/session.py
Original file line number Diff line number Diff line change
Expand Up @@ -167,11 +167,14 @@ def create(self):
self._labels,
observability_options=observability_options,
metadata=metadata,
), MetricsCapture():
) as span, MetricsCapture():
session_pb = api.create_session(
request=request,
metadata=self._database.metadata_with_request_id(
self._database._next_nth_request, 1, metadata
self._database._next_nth_request,
1,
metadata,
span,
),
)
self._session_id = session_pb.name.split("/")[-1]
Expand Down Expand Up @@ -218,7 +221,10 @@ def exists(self):
api.get_session(
name=self.name,
metadata=database.metadata_with_request_id(
database._next_nth_request, 1, metadata
database._next_nth_request,
1,
metadata,
span,
),
)
if span:
Expand Down Expand Up @@ -263,11 +269,14 @@ def delete(self):
},
observability_options=observability_options,
metadata=metadata,
), MetricsCapture():
) as span, MetricsCapture():
api.delete_session(
name=self.name,
metadata=database.metadata_with_request_id(
database._next_nth_request, 1, metadata
database._next_nth_request,
1,
metadata,
span,
),
)

Expand Down
42 changes: 30 additions & 12 deletions google/cloud/spanner_v1/snapshot.py
Original file line number Diff line number Diff line change
Expand Up @@ -104,11 +104,14 @@ def _restart_on_unavailable(
attributes,
observability_options=observability_options,
metadata=metadata,
), MetricsCapture():
) as span, MetricsCapture():
iterator = method(
request=request,
metadata=request_id_manager.metadata_with_request_id(
nth_request, attempt, metadata
nth_request,
attempt,
metadata,
span,
),
)
for item in iterator:
Expand All @@ -133,7 +136,7 @@ def _restart_on_unavailable(
attributes,
observability_options=observability_options,
metadata=metadata,
), MetricsCapture():
) as span, MetricsCapture():
request.resume_token = resume_token
if transaction is not None:
transaction_selector = transaction._make_txn_selector()
Expand All @@ -142,7 +145,10 @@ def _restart_on_unavailable(
iterator = method(
request=request,
metadata=request_id_manager.metadata_with_request_id(
nth_request, attempt, metadata
nth_request,
attempt,
metadata,
span,
),
)
continue
Expand All @@ -160,7 +166,7 @@ def _restart_on_unavailable(
attributes,
observability_options=observability_options,
metadata=metadata,
), MetricsCapture():
) as span, MetricsCapture():
request.resume_token = resume_token
if transaction is not None:
transaction_selector = transaction._make_txn_selector()
Expand All @@ -169,7 +175,10 @@ def _restart_on_unavailable(
iterator = method(
request=request,
metadata=request_id_manager.metadata_with_request_id(
nth_request, attempt, metadata
nth_request,
attempt,
metadata,
span,
),
)
continue
Expand Down Expand Up @@ -745,13 +754,16 @@ def partition_read(
extra_attributes=trace_attributes,
observability_options=getattr(database, "observability_options", None),
metadata=metadata,
), MetricsCapture():
) as span, MetricsCapture():
nth_request = getattr(database, "_next_nth_request", 0)
attempt = AtomicCounter()

def attempt_tracking_method():
all_metadata = database.metadata_with_request_id(
nth_request, attempt.increment(), metadata
nth_request,
attempt.increment(),
metadata,
span,
)
method = functools.partial(
api.partition_read,
Expand Down Expand Up @@ -858,13 +870,16 @@ def partition_query(
trace_attributes,
observability_options=getattr(database, "observability_options", None),
metadata=metadata,
), MetricsCapture():
) as span, MetricsCapture():
nth_request = getattr(database, "_next_nth_request", 0)
attempt = AtomicCounter()

def attempt_tracking_method():
all_metadata = database.metadata_with_request_id(
nth_request, attempt.increment(), metadata
nth_request,
attempt.increment(),
metadata,
span,
)
method = functools.partial(
api.partition_query,
Expand Down Expand Up @@ -1014,13 +1029,16 @@ def begin(self):
self._session,
observability_options=getattr(database, "observability_options", None),
metadata=metadata,
), MetricsCapture():
) as span, MetricsCapture():
nth_request = getattr(database, "_next_nth_request", 0)
attempt = AtomicCounter()

def attempt_tracking_method():
all_metadata = database.metadata_with_request_id(
nth_request, attempt.increment(), metadata
nth_request,
attempt.increment(),
metadata,
span,
)
method = functools.partial(
api.begin_transaction,
Expand Down
16 changes: 1 addition & 15 deletions google/cloud/spanner_v1/testing/interceptors.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@

from grpc_interceptor import ClientInterceptor
from google.api_core.exceptions import Aborted
from google.cloud.spanner_v1.request_id_header import parse_request_id


class MethodCountInterceptor(ClientInterceptor):
Expand Down Expand Up @@ -119,18 +120,3 @@ def stream_request_ids(self):
def reset(self):
self._stream_req_segments.clear()
self._unary_req_segments.clear()


def parse_request_id(request_id_str):
splits = request_id_str.split(".")
version, rand_process_id, client_id, channel_id, nth_request, nth_attempt = list(
map(lambda v: int(v), splits)
)
return (
version,
rand_process_id,
client_id,
channel_id,
nth_request,
nth_attempt,
)
17 changes: 13 additions & 4 deletions google/cloud/spanner_v1/transaction.py
Original file line number Diff line number Diff line change
Expand Up @@ -191,7 +191,10 @@ def wrapped_method(*args, **kwargs):
session=self._session.name,
options=txn_options,
metadata=database.metadata_with_request_id(
nth_request, attempt.increment(), metadata
nth_request,
attempt.increment(),
metadata,
span,
),
)
return method(*args, **kwargs)
Expand Down Expand Up @@ -232,7 +235,7 @@ def rollback(self):
self._session,
observability_options=observability_options,
metadata=metadata,
), MetricsCapture():
) as span, MetricsCapture():
attempt = AtomicCounter(0)
nth_request = database._next_nth_request

Expand All @@ -243,7 +246,10 @@ def wrapped_method(*args, **kwargs):
session=self._session.name,
transaction_id=self._transaction_id,
metadata=database.metadata_with_request_id(
nth_request, attempt.value, metadata
nth_request,
attempt.value,
metadata,
span,
),
)
return method(*args, **kwargs)
Expand Down Expand Up @@ -334,7 +340,10 @@ def wrapped_method(*args, **kwargs):
api.commit,
request=request,
metadata=database.metadata_with_request_id(
nth_request, attempt.value, metadata
nth_request,
attempt.value,
metadata,
span,
),
)
return method(*args, **kwargs)
Expand Down
Loading
Loading