diff --git a/CHANGELOG.md b/CHANGELOG.md index 04c5304a..f5f6678e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,6 @@ +## 3.32.1 ## +* Add the `ydb.topic.reader.received.messages` counter and optional topic reader names for distinguishing reader metric series; bump the metrics build-info token to `ydb-sdk-metrics/0.2.0` + ## 3.32.0 ## * Add `TableClient.read_rows` (sync and async) to read rows by primary key without a transaction * Add the `ydb.query.session.closed` counter for query session pool closures, labeled by pool name and a standardized closure reason; metrics-enabled clients now advertise `ydb-sdk-metrics/0.2.0` in `x-ydb-sdk-build-info` diff --git a/docs/observability.rst b/docs/observability.rst index f332cf7f..b2ec2e7c 100644 --- a/docs/observability.rst +++ b/docs/observability.rst @@ -352,6 +352,9 @@ adapter maps them to instruments on the ``"ydb.sdk"`` meter): * - ``ydb.client.retry.attempts`` - Histogram - Number of attempts performed for one logical retried operation. + * - ``ydb.topic.reader.received.messages`` + - Counter (``{message}``) + - Messages accepted into the local SDK topic reader buffer. Attributes ~~~~~~~~~~ @@ -402,6 +405,15 @@ initial attach handshake does not count as closing an active session, and standa ``QuerySession`` instances do not publish pool metrics. A session publishes at most one closure event; the first terminal reason wins. +``ydb.topic.reader.received.messages`` carries ``endpoint``, ``database``, ``topic``, +``consumer``, and ``reader.name``. For reads without a consumer, ``consumer`` is an +empty string. The user can pass ``reader_name`` to +``TopicClient.reader``; otherwise the SDK generates a process-local ``reader-N`` value +once for the logical reader and preserves it across reconnects. Use a rate function on +this cumulative counter to diagnose incoming progress. A gap between received and an +application-level delivered-message metric can indicate that the application is not +consuming data or that decoding is failing. + Writing a Custom Metrics Backend ~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ diff --git a/docs/topic.rst b/docs/topic.rst index bd545569..9476bbf4 100644 --- a/docs/topic.rst +++ b/docs/topic.rst @@ -378,8 +378,14 @@ Reader Parameters consumer="my-consumer", buffer_size_bytes=50 * 1024 * 1024, # client-side buffer (default: 50 MB) buffer_release_threshold=0.5, # see below (default: 0.5) + reader_name="payments-worker", # optional name used in reader metrics ) +``reader_name`` is an optional stable name for distinguishing topic readers in +observability metrics. If it is omitted or empty, the SDK assigns a process-local name +in the form ``reader-N``. Explicit names are not required to be unique: readers using +the same name contribute to the same metric series when their other attributes match. + ``buffer_size_bytes`` controls how many bytes the server is allowed to send before the client signals that it is ready for more. The server will not exceed this limit. diff --git a/examples/opentelemetry/README.md b/examples/opentelemetry/README.md index 33fd6899..d18b00f7 100644 --- a/examples/opentelemetry/README.md +++ b/examples/opentelemetry/README.md @@ -51,8 +51,11 @@ Grafana is provisioned with the **YDB Python SDK Metrics** dashboard. It uses Prometheus queries for SDK metrics such as `db_client_operation_duration`, `ydb_client_operation_failed`, `ydb_query_session_count`, `ydb_query_session_pending_requests`, `ydb_query_session_create_time`, and -`ydb_client_retry_duration`. Use Grafana Explore for ad-hoc traces through Tempo -and metrics through Prometheus. +`ydb_client_retry_duration`. Topic readers also export the cumulative +`ydb_topic_reader_received_messages_total` counter; use `rate(...)` and group by +`topic`, `consumer`, or `reader_name` to inspect incoming progress (the OpenTelemetry +attribute `reader.name` is normalized to `reader_name` by the Prometheus exporter). Use +Grafana Explore for ad-hoc traces through Tempo and metrics through Prometheus. The SDK configures explicit OpenTelemetry histogram bucket boundaries for its own duration and retry-attempt metrics. Duration values are recorded in seconds, diff --git a/tests/observability/test_metrics.py b/tests/observability/test_metrics.py index 94c8ead6..6a13031f 100644 --- a/tests/observability/test_metrics.py +++ b/tests/observability/test_metrics.py @@ -103,10 +103,12 @@ def test_metrics_registry_records_all_instruments(metrics_setup, monkeypatch): QUERY_SESSION_TIMEOUTS, RETRY_ATTEMPTS, RETRY_DURATION, + TOPIC_READER_RECEIVED_MESSAGES, ATTEMPT_BUCKETS, DURATION_BUCKETS_SECONDS, RETRY_DURATION_BUCKETS_SECONDS, SessionMetrics, + TopicReaderMetrics, create_metrics_operation, record_query_session_count, record_query_session_create_time, @@ -133,6 +135,11 @@ def test_metrics_registry_records_all_instruments(metrics_setup, monkeypatch): record_query_session_pending_requests(1, "main") record_query_session_timeout("main") record_retry_metrics(0.75, 3) + TopicReaderMetrics( + object(), + consumer_name=None, + reader_name="reader-test", + ).record_received_messages(1, "/Root/events") metrics = _metrics_by_name(metrics_setup) @@ -148,6 +155,7 @@ def test_metrics_registry_records_all_instruments(metrics_setup, monkeypatch): QUERY_SESSION_TIMEOUTS, RETRY_ATTEMPTS, RETRY_DURATION, + TOPIC_READER_RECEIVED_MESSAGES, } assert metrics[CLIENT_OPERATION_DURATION].unit == "s" assert metrics[CLIENT_OPERATION_FAILED].unit == "{command}" @@ -169,6 +177,7 @@ def test_metrics_registry_records_all_instruments(metrics_setup, monkeypatch): ) assert _single_point_from_metrics(metrics, RETRY_DURATION).explicit_bounds == RETRY_DURATION_BUCKETS_SECONDS assert _single_point_from_metrics(metrics, RETRY_ATTEMPTS).explicit_bounds == ATTEMPT_BUCKETS + assert metrics[TOPIC_READER_RECEIVED_MESSAGES].unit == "{message}" def test_metrics_registry_supports_old_histogram_api(): @@ -294,6 +303,12 @@ def test_metrics_registry_is_noop_without_meter(monkeypatch): pool_metrics.close() +def test_metrics_build_info_token_version(metrics_setup): + from ydb.observability import sdk_build_info_tokens + + assert sdk_build_info_tokens() == ["ydb-sdk-metrics/0.2.0"] + + def test_metrics_operation_records_duration_once(metrics_setup, monkeypatch): from ydb.observability.metrics import CLIENT_OPERATION_DURATION, create_metrics_operation @@ -1698,3 +1713,124 @@ def fake_span_ctx(**kwargs): assert qs._session_metrics._counted assert _single_point_for_pool(metrics_setup, QUERY_SESSION_COUNT, "async-session-pool").value == 1 qs._close_session() + + +def test_topic_reader_received_messages_attributes(metrics_setup): + from tests.observability.conftest import FakeDriverConfig + from ydb.observability.metrics import ( + TOPIC_READER_RECEIVED_MESSAGES, + TopicReaderMetrics, + ) + + class FakeDriver: + _driver_config = FakeDriverConfig( + endpoint="grpc://localhost:2136", + database="/Root", + ) + + reader_metrics = TopicReaderMetrics( + FakeDriver(), + consumer_name="analytics", + reader_name="payments-worker", + ) + + reader_metrics.record_received_messages( + count=2, + topic="/Root/events", + ) + reader_metrics.record_received_messages( + count=3, + topic="/Root/events", + ) + + point = _single_point( + metrics_setup, + TOPIC_READER_RECEIVED_MESSAGES, + ) + + assert point.value == 5 + assert point.attributes == { + "endpoint": "localhost:2136", + "database": "/Root", + "topic": "/Root/events", + "consumer": "analytics", + "reader.name": "payments-worker", + } + + +def test_topic_reader_received_messages_without_consumer(metrics_setup): + from ydb.observability.metrics import ( + TOPIC_READER_RECEIVED_MESSAGES, + TopicReaderMetrics, + ) + + reader_metrics = TopicReaderMetrics( + object(), + consumer_name=None, + reader_name="reader-42", + ) + + reader_metrics.record_received_messages( + count=1, + topic="/Root/events", + ) + + point = _single_point( + metrics_setup, + TOPIC_READER_RECEIVED_MESSAGES, + ) + + assert point.attributes == { + "endpoint": "", + "database": "", + "topic": "/Root/events", + "consumer": "", + "reader.name": "reader-42", + } + + +def test_topic_reader_received_messages_separates_topics(metrics_setup): + from ydb.observability.metrics import ( + TOPIC_READER_RECEIVED_MESSAGES, + TopicReaderMetrics, + ) + + reader_metrics = TopicReaderMetrics( + object(), + consumer_name="analytics", + reader_name="worker", + ) + + reader_metrics.record_received_messages(2, "/Root/a") + reader_metrics.record_received_messages(3, "/Root/b") + + values = { + point.attributes["topic"]: point.value + for point in _points( + metrics_setup, + TOPIC_READER_RECEIVED_MESSAGES, + ) + } + + assert values == { + "/Root/a": 2, + "/Root/b": 3, + } + + +def test_topic_reader_received_messages_ignores_nonpositive_values(metrics_setup): + from ydb.observability.metrics import ( + TOPIC_READER_RECEIVED_MESSAGES, + TopicReaderMetrics, + ) + + reader_metrics = TopicReaderMetrics( + object(), + consumer_name="analytics", + reader_name="worker", + ) + + reader_metrics.record_received_messages(0, "/Root/events") + reader_metrics.record_received_messages(-1, "/Root/events") + + assert TOPIC_READER_RECEIVED_MESSAGES not in _metrics_by_name(metrics_setup) diff --git a/ydb/_grpc/grpcwrapper/ydb_topic.py b/ydb/_grpc/grpcwrapper/ydb_topic.py index 19b38bf0..0de3a749 100644 --- a/ydb/_grpc/grpcwrapper/ydb_topic.py +++ b/ydb/_grpc/grpcwrapper/ydb_topic.py @@ -486,11 +486,14 @@ class InitRequest(IToProto): topics_read_settings: List["StreamReadMessage.InitRequest.TopicReadSettings"] consumer: Optional[str] auto_partitioning_support: bool + reader_name: Optional[str] = None def to_proto(self) -> ydb_topic_pb2.StreamReadMessage.InitRequest: res = ydb_topic_pb2.StreamReadMessage.InitRequest() if self.consumer is not None: res.consumer = self.consumer + if self.reader_name is not None: + res.reader_name = self.reader_name for settings in self.topics_read_settings: res.topics_read_settings.append(settings.to_proto()) res.auto_partitioning_support = self.auto_partitioning_support diff --git a/ydb/_grpc/grpcwrapper/ydb_topic_test.py b/ydb/_grpc/grpcwrapper/ydb_topic_test.py index b9e30603..4c79f2d1 100644 --- a/ydb/_grpc/grpcwrapper/ydb_topic_test.py +++ b/ydb/_grpc/grpcwrapper/ydb_topic_test.py @@ -2,7 +2,7 @@ from google.protobuf.json_format import MessageToDict -from ydb._grpc.grpcwrapper.ydb_topic import OffsetsRange +from ydb._grpc.grpcwrapper.ydb_topic import OffsetsRange, StreamReadMessage from .ydb_topic import AlterTopicRequest from .ydb_topic_public_types import ( AlterTopicRequestParams, @@ -96,3 +96,27 @@ def test_alter_topic_request_from_public_to_proto(): } assert msg_dict == expected_dict + + +def test_stream_read_init_request_serializes_reader_name(): + request = StreamReadMessage.InitRequest( + topics_read_settings=[], + consumer="analytics", + auto_partitioning_support=True, + reader_name="payments-worker", + ) + + proto = request.to_proto() + + assert proto.reader_name == "payments-worker" + + +def test_stream_read_init_request_omits_reader_name_by_default(): + request = StreamReadMessage.InitRequest( + topics_read_settings=[], + consumer="analytics", + auto_partitioning_support=True, + ) + + assert request.reader_name is None + assert request.to_proto().reader_name == "" diff --git a/ydb/_topic_reader/topic_reader.py b/ydb/_topic_reader/topic_reader.py index e4545f3a..9123b8a7 100644 --- a/ydb/_topic_reader/topic_reader.py +++ b/ydb/_topic_reader/topic_reader.py @@ -60,7 +60,15 @@ class PublicReaderSettings: buffer_release_threshold: float = 0.5 """Min fraction of buffer_size_bytes to accumulate before sending a new ReadRequest (0.0 = immediately after every batch).""" + reader_name: Optional[str] = None + """Optional stable reader name used to distinguish reader metric series.""" + def __post_init__(self): + if self.reader_name is not None and not isinstance( + self.reader_name, + str, + ): + raise TypeError("Unsupported type for reader_name field: '%s'" % type(self.reader_name)) if not (0.0 <= self.buffer_release_threshold <= 1.0): raise ValueError("buffer_release_threshold must be in [0.0, 1.0], got %s" % self.buffer_release_threshold) # check possible create init message @@ -87,6 +95,7 @@ def _init_message(self) -> StreamReadMessage.InitRequest: topics_read_settings=list(map(PublicTopicSelector._to_topic_read_settings, selectors)), # type: ignore consumer=self.consumer, auto_partitioning_support=self.auto_partitioning_support, + reader_name=self.reader_name, ) def _retry_settings(self) -> RetrySettings: diff --git a/ydb/_topic_reader/topic_reader_asyncio.py b/ydb/_topic_reader/topic_reader_asyncio.py index e5c9f3a5..15d10b50 100644 --- a/ydb/_topic_reader/topic_reader_asyncio.py +++ b/ydb/_topic_reader/topic_reader_asyncio.py @@ -36,6 +36,8 @@ from ..query.base import TxEvent +from ..observability.metrics import TopicReaderMetrics + if typing.TYPE_CHECKING: from ..query.transaction import BaseQueryTxContext @@ -235,6 +237,8 @@ class ReaderReconnector: _settings: topic_reader.PublicReaderSettings _driver: Driver _background_tasks: Set[Task] + _reader_name: str + _metrics: TopicReaderMetrics _state_changed: asyncio.Event _stream_reader: Optional["ReaderStream"] @@ -251,6 +255,12 @@ def __init__( self._id = ReaderReconnector._static_reader_reconnector_counter.inc_and_get() self._settings = settings self._driver = driver + self._reader_name = settings.reader_name or "reader-%d" % self._id + self._metrics = TopicReaderMetrics( + driver, + consumer_name=settings.consumer, + reader_name=self._reader_name, + ) self._loop = loop if loop is not None else asyncio.get_running_loop() self._background_tasks = set() logger.debug("init reader reconnector id=%s", self._id) @@ -270,7 +280,12 @@ async def _connection_loop(self): return try: logger.debug("reader %s connect attempt %s", self._id, attempt) - self._stream_reader = await ReaderStream.create(self._id, self._driver, self._settings) + self._stream_reader = await ReaderStream.create( + self._id, + self._driver, + self._settings, + metrics=self._metrics, + ) logger.debug("reader %s connected stream %s", self._id, self._stream_reader._id) attempt = 0 self._state_changed.set() @@ -516,6 +531,7 @@ class ReaderStream: _pending_buffer_release_bytes: int _decode_executor: Optional[concurrent.futures.Executor] _decoders: Dict[int, typing.Callable[[bytes], bytes]] # dict[codec_code] func(encoded_bytes)->decoded_bytes + _metrics: Optional[TopicReaderMetrics] if typing.TYPE_CHECKING: _batches_to_decode: asyncio.Queue[datatypes.PublicBatch] @@ -537,6 +553,7 @@ def __init__( reader_reconnector_id: int, settings: topic_reader.PublicReaderSettings, get_token_function: Optional[Callable[[], str]] = None, + metrics: Optional[TopicReaderMetrics] = None, ): self._loop = asyncio.get_running_loop() self._id = ReaderStream._static_id_counter.inc_and_get() @@ -572,6 +589,8 @@ def __init__( self._settings = settings + self._metrics = metrics + logger.debug("created ReaderStream id=%s reconnector=%s", self._id, self._reader_reconnector_id) @staticmethod @@ -579,6 +598,8 @@ async def create( reader_reconnector_id: int, driver: SupportedDriverType, settings: topic_reader.PublicReaderSettings, + *, + metrics: TopicReaderMetrics, ) -> "ReaderStream": stream = GrpcWrapperAsyncIO(StreamReadMessage.FromServer.from_proto) reader = None @@ -590,6 +611,7 @@ async def create( reader_reconnector_id, settings, get_token_function=creds.get_auth_token if creds else None, + metrics=metrics, ) await reader._start(stream, settings._init_message()) except BaseException: @@ -927,6 +949,11 @@ def _on_read_response(self, message: StreamReadMessage.ReadResponse): batches = self._read_response_to_batches(message) for batch in batches: self._batches_to_decode.put_nowait(batch) + if self._metrics is not None: + self._metrics.record_received_messages( + count=len(batch.messages), + topic=batch._partition_session.topic_path, + ) def _on_commit_response(self, message: StreamReadMessage.CommitOffsetResponse): for partition_offset in message.partitions_committed_offsets: diff --git a/ydb/_topic_reader/topic_reader_asyncio_test.py b/ydb/_topic_reader/topic_reader_asyncio_test.py index c018604f..938c211e 100644 --- a/ydb/_topic_reader/topic_reader_asyncio_test.py +++ b/ydb/_topic_reader/topic_reader_asyncio_test.py @@ -29,6 +29,7 @@ wait_for_fast, WaitConditionError, ) +from ..observability.metrics import TopicReaderMetrics # Workaround for good IDE and universal for runtime if typing.TYPE_CHECKING: @@ -282,6 +283,51 @@ def batch_size(): else: await wait_condition(lambda: batch_size() > initial_batch_size) + async def test_received_messages_metric_is_recorded_before_decoding( + self, + stream, + default_reader_settings, + ): + metrics = mock.Mock(spec=TopicReaderMetrics) + reader = await self.get_started_reader( + stream, + default_reader_settings, + metrics=metrics, + ) + partition_session = datatypes.PartitionSession( + id=self.partition_session_id, + topic_path=default_reader_settings.topic, + partition_id=4, + state=datatypes.PartitionSession.State.Active, + committed_offset=self.partition_session_committed_offset, + reader_reconnector_id=self.default_reader_reconnector_id, + reader_stream_id=reader._id, + ) + reader._partition_sessions[partition_session.id] = partition_session + + def assert_batch_is_accepted_for_decoding(*, count, topic): + assert reader._batches_to_decode.qsize() == 1 + assert partition_session.id not in reader._message_batches + assert topic == default_reader_settings.topic + + metrics.record_received_messages.side_effect = assert_batch_is_accepted_for_decoding + + try: + await self.send_batch( + reader, + [ + self.create_message(partition_session, 1, 1), + self.create_message(partition_session, 2, 1), + ], + ) + + metrics.record_received_messages.assert_called_once_with( + count=2, + topic=default_reader_settings.topic, + ) + finally: + await reader.close(False) + async def test_unknown_error(self, stream, stream_reader_finish_with_error): class TestError(Exception): pass @@ -1620,8 +1666,11 @@ async def stream_create( reader_reconnector_id: int, driver: SupportedDriverType, settings: PublicReaderSettings, + *, + metrics: TopicReaderMetrics, ): nonlocal stream_index + metric_contexts.append(metrics) stream_index += 1 if stream_index == 1: return reader_stream_mock_with_error @@ -1630,9 +1679,16 @@ async def stream_create( else: raise Exception("unexpected create stream") + metric_contexts = [] + driver = mock.Mock() + driver._driver_config = None with mock.patch.object(ReaderStream, "create", stream_create): - reconnector = ReaderReconnector(mock.Mock(), PublicReaderSettings("", "")) + reconnector = ReaderReconnector(driver, PublicReaderSettings("", "")) await wait_for_fast(reconnector.wait_message()) + assert len(metric_contexts) == 2 + assert metric_contexts[0] is metric_contexts[1] is reconnector._metrics + assert reconnector._metrics._base_attributes["reader.name"] == "reader-%d" % reconnector._id + await reconnector.close(flush=False) reader_stream_mock_with_error.wait_error.assert_any_await() reader_stream_mock_with_error.wait_messages.assert_any_await() @@ -1668,19 +1724,53 @@ async def wait_forever(): create_calls = 0 - async def stream_create(reader_reconnector_id, driver, settings): + async def stream_create(reader_reconnector_id, driver, settings, *, metrics): nonlocal create_calls create_calls += 1 return stream1 if create_calls == 1 else stream2 + driver = mock.Mock() + driver._driver_config = None with mock.patch.object(ReaderStream, "create", stream_create): - reconnector = ReaderReconnector(mock.Mock(), PublicReaderSettings("", "")) + reconnector = ReaderReconnector(driver, PublicReaderSettings("", "")) await asyncio.wait_for(finally_close_started.wait(), timeout=2) await asyncio.wait_for(reconnector.close(flush=False), timeout=5) # The loop stopped on close instead of reconnecting into a second (zombie) stream. assert create_calls == 1 + async def test_reader_name_uses_user_value_or_process_local_sequence(self): + async def stream_create(reader_reconnector_id, driver, settings, *, metrics): + await asyncio.Future() + + driver = mock.Mock() + driver._driver_config = None + with mock.patch.object(ReaderStream, "create", stream_create): + generated_from_none = ReaderReconnector( + driver, + PublicReaderSettings("consumer", "topic"), + ) + generated_from_empty = ReaderReconnector( + driver, + PublicReaderSettings("consumer", "topic", reader_name=""), + ) + named = ReaderReconnector( + driver, + PublicReaderSettings("consumer", "topic", reader_name="payments-worker"), + ) + + assert generated_from_none._reader_name == "reader-%d" % generated_from_none._id + assert generated_from_empty._reader_name == "reader-%d" % generated_from_empty._id + assert generated_from_none._reader_name != generated_from_empty._reader_name + assert named._reader_name == "payments-worker" + assert named._metrics._base_attributes["reader.name"] == "payments-worker" + + await asyncio.gather( + generated_from_none.close(flush=False), + generated_from_empty.close(flush=False), + named.close(flush=False), + ) + async def test_create_closes_inflight_stream_on_cancel(self, default_reader_settings): # If create() is cancelled (e.g. reader.close() cancels the connection loop during a # reconnect) while parked on the init handshake, the in-flight gRPC stream must be @@ -1698,11 +1788,19 @@ async def start(self, driver, stub, method): driver = mock.Mock() driver._credentials = None + driver._driver_config = None with mock.patch.object(topic_reader_asyncio, "GrpcWrapperAsyncIO", FakeStream): # Real create(); no InitResponse is sent, so it parks inside _start() on # `await stream.receive()` (the only reachable cancellation point in create()). - create_task = asyncio.create_task(ReaderStream.create(7, driver, default_reader_settings)) + create_task = asyncio.create_task( + ReaderStream.create( + 7, + driver, + default_reader_settings, + metrics=mock.Mock(spec=TopicReaderMetrics), + ) + ) await wait_condition(lambda: bool(built) and not built[0].from_client.empty()) assert not create_task.done() @@ -1930,3 +2028,27 @@ async def test_threshold_one_flushes_when_bytes_match_buffer_size(self, stream, assert msg.client_message.bytes_size == 1000 await reader.close(False) + + +def test_reader_settings_forward_reader_name(): + settings = PublicReaderSettings( + consumer="analytics", + topic="/Root/events", + reader_name="payments-worker", + ) + + init_message = settings._init_message() + + assert init_message.reader_name == "payments-worker" + assert init_message.to_proto().reader_name == "payments-worker" + + +def test_reader_settings_positional_buffer_size_is_preserved(): + settings = PublicReaderSettings( + "analytics", + "/Root/events", + 1024, + ) + + assert settings.buffer_size_bytes == 1024 + assert settings.reader_name is None diff --git a/ydb/observability/metrics.py b/ydb/observability/metrics.py index 63cc7008..89219cd8 100644 --- a/ydb/observability/metrics.py +++ b/ydb/observability/metrics.py @@ -38,6 +38,7 @@ QUERY_SESSION_MIN = "ydb.query.session.min" RETRY_ATTEMPTS = "ydb.client.retry.attempts" RETRY_DURATION = "ydb.client.retry.duration" +TOPIC_READER_RECEIVED_MESSAGES = "ydb.topic.reader.received.messages" METRICS_SDK_BUILD_INFO = "ydb-sdk-metrics/0.2.0" @@ -496,6 +497,48 @@ def __exit__(self, exc_type, exc_val, exc_tb): _NOOP_CM = _NoopContext() +class TopicReaderMetrics: + """Metric context shared by all streams of one logical topic reader.""" + + __slots__ = ("_base_attributes",) + + def __init__( + self, + driver, + consumer_name: Optional[str], + reader_name: str, + ) -> None: + driver_config = getattr( + driver, + "_driver_config", + None, + ) + self._base_attributes = _build_ydb_metrics_attrs(driver_config) + self._base_attributes.update( + { + "consumer": consumer_name or "", + "reader.name": reader_name, + } + ) + + def record_received_messages( + self, + count: int, + topic: str, + ) -> None: + if count <= 0 or not is_metrics_enabled(): + return + + attributes = dict(self._base_attributes) + attributes["topic"] = topic + + _provider.add( + TOPIC_READER_RECEIVED_MESSAGES, + count, + attributes, + ) + + class SessionMetrics: """Per-session query-session-count bookkeeping, kept out of the session's own code. diff --git a/ydb/opentelemetry/metrics_plugin.py b/ydb/opentelemetry/metrics_plugin.py index 3eed89df..06c52e98 100644 --- a/ydb/opentelemetry/metrics_plugin.py +++ b/ydb/opentelemetry/metrics_plugin.py @@ -34,6 +34,7 @@ DURATION_BUCKETS_SECONDS, GaugeCallback, RETRY_DURATION_BUCKETS_SECONDS, + TOPIC_READER_RECEIVED_MESSAGES, _get_metrics_provider, ) @@ -114,6 +115,14 @@ def __init__(self, meter: Meter) -> None: unit="{request}", description="Number of requests waiting for a YDB query session.", ), + TOPIC_READER_RECEIVED_MESSAGES: meter.create_counter( + TOPIC_READER_RECEIVED_MESSAGES, + unit="{message}", + description=( + "Number of messages accepted into the local SDK topic reader buffer " + "for an active partition session." + ), + ), } def record(self, name: str, value: float, attributes: Optional[Dict[str, Any]] = None) -> None: diff --git a/ydb/topic.py b/ydb/topic.py index 98859293..0c88ebe5 100644 --- a/ydb/topic.py +++ b/ydb/topic.py @@ -295,6 +295,7 @@ def reader( auto_partitioning_support: Optional[bool] = True, # Auto partitioning feature flag. Default - True. event_handler: Optional[TopicReaderEvents.EventHandler] = None, buffer_release_threshold: float = 0.5, + reader_name: Optional[str] = None, ) -> TopicReaderAsyncIO: logger.debug("Create reader for topic=%s consumer=%s", topic, consumer) @@ -631,6 +632,7 @@ def reader( auto_partitioning_support: Optional[bool] = True, # Auto partitioning feature flag. Default - True. event_handler: Optional[TopicReaderEvents.EventHandler] = None, buffer_release_threshold: float = 0.5, + reader_name: Optional[str] = None, ) -> TopicReader: logger.debug("Create reader for topic=%s consumer=%s", topic, consumer) if not decoder_executor: