Skip to content
Open
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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
* Add optional topic reader names for identifying logical readers in server diagnostics and SDK logs
* Deprecated the table client scan query methods — `TableClient.scan_query`, `TableClient.async_scan_query` and the async `ydb.aio.TableClient.scan_query`: they now emit a `DeprecationWarning` and keep working as before, use QueryService (`ydb.QuerySessionPool` / `ydb.aio.QuerySessionPool`) instead
* Mark the package as typed so type checkers use the SDK's inline annotations

Expand Down
5 changes: 5 additions & 0 deletions docs/topic.rst
Original file line number Diff line number Diff line change
Expand Up @@ -378,8 +378,13 @@ 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 diagnostic name
)

``reader_name`` identifies a logical reader in server-side diagnostics and SDK logs. If it is
omitted or empty, the SDK generates a process-local name in the form ``reader-N``. The name
remains unchanged when the reader reconnects. Explicit names do not have to be unique.

``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.

Expand Down
6 changes: 5 additions & 1 deletion examples/topic/reader_async_example.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,11 @@ async def connect():
connection_string="grpc://localhost:2135?database=/local",
credentials=ydb.credentials.AnonymousCredentials(),
)
reader = db.topic_client.reader("/local/topic", consumer="consumer")
reader = db.topic_client.reader(
"/local/topic",
consumer="consumer",
reader_name="payments-worker",
)
return reader


Expand Down
6 changes: 5 additions & 1 deletion examples/topic/reader_example.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,11 @@ def connect():
connection_string="grpc://localhost:2135?database=/local",
credentials=ydb.credentials.AnonymousCredentials(),
)
reader = db.topic_client.reader("/local/topic", consumer="consumer")
reader = db.topic_client.reader(
"/local/topic",
consumer="consumer",
reader_name="payments-worker",
)
return reader


Expand Down
3 changes: 3 additions & 0 deletions ydb/_grpc/grpcwrapper/ydb_topic.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
24 changes: 23 additions & 1 deletion ydb/_grpc/grpcwrapper/ydb_topic_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -96,3 +96,25 @@ 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",
)

assert request.to_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 == ""
8 changes: 7 additions & 1 deletion ydb/_topic_reader/topic_reader.py
Original file line number Diff line number Diff line change
Expand Up @@ -60,13 +60,18 @@ 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 reader name used to identify this logical reader in diagnostics."""

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
_ = self._init_message()

def _init_message(self) -> StreamReadMessage.InitRequest:
def _init_message(self, *, reader_name: Optional[str] = None) -> StreamReadMessage.InitRequest:
if self.consumer is not None and not isinstance(self.consumer, str):
raise TypeError("Unsupported type for customer field: '%s'" % type(self.consumer))

Expand All @@ -87,6 +92,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 if reader_name is None else reader_name,
)

def _retry_settings(self) -> RetrySettings:
Expand Down
91 changes: 64 additions & 27 deletions ydb/_topic_reader/topic_reader_asyncio.py
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ def __init__(self):
class PublicAsyncIOReader:
_loop: asyncio.AbstractEventLoop
_closed: bool
_log_prefix: str
_settings: topic_reader.PublicReaderSettings
_reconnector: ReaderReconnector
_parent: typing.Any # need for prevent close parent client by GC
Expand All @@ -100,8 +101,10 @@ def __init__(
):
self._loop = asyncio.get_running_loop()
self._closed = False
self._log_prefix = "topic reader"
self._settings = settings
self._reconnector = ReaderReconnector(driver, settings, self._loop)
self._log_prefix = self._reconnector._log_prefix
self._parent = _parent

async def __aenter__(self):
Expand All @@ -111,13 +114,19 @@ async def __aexit__(self, exc_type, exc_val, exc_tb):
await self.close()

def __del__(self):
if not self._closed:
try:
logger.debug("Topic reader was not closed properly. Consider using method close().")
task = self._loop.create_task(self.close(flush=False))
task.set_name("close reader")
except BaseException:
logger.warning("Something went wrong during reader close in __del__")
if getattr(self, "_closed", True):
return

loop = getattr(self, "_loop", None)
if getattr(self, "_reconnector", None) is None or loop is None or loop.is_closed() or not loop.is_running():
return

try:
logger.debug("%s was not closed properly. Consider using method close().", self._log_prefix)
task = loop.create_task(self.close(flush=False))
task.set_name("close reader")
except BaseException:
logger.warning("%s failed to close in __del__", self._log_prefix)

async def wait_message(self):
"""
Expand All @@ -140,7 +149,7 @@ async def receive_batch(

use asyncio.wait_for for wait with timeout.
"""
logger.debug("receive_batch max_messages=%s max_bytes=%s", max_messages, max_bytes)
logger.debug("%s receive_batch max_messages=%s max_bytes=%s", self._log_prefix, max_messages, max_bytes)
await self._reconnector.wait_message()
return self._reconnector.receive_batch_nowait(
max_messages=max_messages,
Expand All @@ -163,7 +172,13 @@ async def receive_batch_with_tx(

use asyncio.wait_for for wait with timeout.
"""
logger.debug("receive_batch_with_tx tx=%s max_messages=%s max_bytes=%s", tx, max_messages, max_bytes)
logger.debug(
"%s receive_batch_with_tx tx=%s max_messages=%s max_bytes=%s",
self._log_prefix,
tx,
max_messages,
max_bytes,
)
await self._reconnector.wait_message()
return self._reconnector.receive_batch_with_tx_nowait(
tx=tx,
Expand All @@ -177,7 +192,7 @@ async def receive_message(self) -> typing.Optional[datatypes.PublicMessage]:

use asyncio.wait_for for wait with timeout.
"""
logger.debug("receive_message")
logger.debug("%s receive_message", self._log_prefix)
await self._reconnector.wait_message()
return self._reconnector.receive_message_nowait()

Expand All @@ -188,7 +203,7 @@ def commit(self, batch: typing.Union[datatypes.PublicMessage, datatypes.PublicBa
For the method no way check the commit result
(for example if lost connection - commits will not re-send and committed messages will receive again).
"""
logger.debug("commit message or batch")
logger.debug("%s commit message or batch", self._log_prefix)
if self._settings.consumer is None:
raise issues.Error("Commit operations are not supported for topic reader without consumer.")

Expand All @@ -207,7 +222,7 @@ async def commit_with_ack(self, batch: typing.Union[datatypes.PublicMessage, dat
before receive commit ack. Message may be acked or not (if not - it will send in other read session,
to this or other reader).
"""
logger.debug("commit_with_ack message or batch")
logger.debug("%s commit_with_ack message or batch", self._log_prefix)
if self._settings.consumer is None:
raise issues.Error("Commit operations are not supported for topic reader without consumer.")

Expand All @@ -218,10 +233,10 @@ async def close(self, flush: bool = True):
if self._closed:
raise TopicReaderClosedError()

logger.debug("Close topic reader")
logger.debug("%s close", self._log_prefix)
self._closed = True
await self._reconnector.close(flush)
logger.debug("Topic reader was closed")
logger.debug("%s was closed", self._log_prefix)

@property
def read_session_id(self) -> Optional[str]:
Expand All @@ -235,6 +250,8 @@ class ReaderReconnector:
_settings: topic_reader.PublicReaderSettings
_driver: Driver
_background_tasks: Set[Task]
_reader_name: str
_log_prefix: str

_state_changed: asyncio.Event
_stream_reader: Optional["ReaderStream"]
Expand All @@ -251,9 +268,11 @@ 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._log_prefix = "topic reader reader_name=%r reader_id=%s" % (self._reader_name, self._id)
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)
logger.debug("%s initialize reconnector", self._log_prefix)

self._state_changed = asyncio.Event()
self._stream_reader = None
Expand All @@ -269,21 +288,30 @@ async def _connection_loop(self):
if self._closed:
return
try:
logger.debug("reader %s connect attempt %s", self._id, attempt)
self._stream_reader = await ReaderStream.create(self._id, self._driver, self._settings)
logger.debug("reader %s connected stream %s", self._id, self._stream_reader._id)
logger.debug("%s connect attempt=%s", self._log_prefix, attempt)
self._stream_reader = await ReaderStream.create(
self._id,
self._driver,
self._settings,
reader_name=self._reader_name,
)
logger.debug("%s connected stream_id=%s", self._log_prefix, self._stream_reader._id)
attempt = 0
self._state_changed.set()
await self._stream_reader.wait_error()
except BaseException as err:
logger.debug("reader %s, attempt %s connection loop error %s", self._id, attempt, err)
logger.debug("%s connection attempt=%s failed: %s", self._log_prefix, attempt, err)
retry_info = check_retriable_error(err, self._settings._retry_settings(), attempt)
if not retry_info.is_retriable:
logger.debug("reader %s stop connection loop due to %s", self._id, err)
logger.debug("%s stop connection loop: %s", self._log_prefix, err)
self._set_first_error(err)
return

logger.debug("sleep before retry for %s seconds", retry_info.sleep_timeout_seconds)
logger.debug(
"%s sleep before retry for %s seconds",
self._log_prefix,
retry_info.sleep_timeout_seconds,
)

await asyncio.sleep(retry_info.sleep_timeout_seconds)

Expand Down Expand Up @@ -465,7 +493,7 @@ def commit(self, batch: datatypes.ICommittable) -> datatypes.PartitionSession.Co
return self._stream_reader.commit(batch)

async def close(self, flush: bool):
logger.debug("reader reconnector %s close", self._id)
logger.debug("%s close reconnector", self._log_prefix)
# Mark closed so the connection loop won't start a new stream, then close the
# current stream with the requested flush before cancelling the loop. On a normal
# close this flushes pending commits; cancelling the loop first would let it close
Expand Down Expand Up @@ -506,6 +534,7 @@ class ReaderStream:
_loop: asyncio.AbstractEventLoop
_id: int
_reader_reconnector_id: int
_reader_name: str
_session_id: str
_stream: Optional[IGrpcWrapperAsyncIO]
_started: bool
Expand Down Expand Up @@ -537,12 +566,16 @@ def __init__(
reader_reconnector_id: int,
settings: topic_reader.PublicReaderSettings,
get_token_function: Optional[Callable[[], str]] = None,
*,
reader_name: str,
):
self._loop = asyncio.get_running_loop()
self._id = ReaderStream._static_id_counter.inc_and_get()
self._reader_reconnector_id = reader_reconnector_id
self._reader_name = reader_name
self._session_id = "not initialized"
self._log_prefix = "reader %s stream %s session=%s" % (
self._log_prefix = "topic reader reader_name=%r reader_id=%s stream_id=%s session_id=%s" % (
self._reader_name,
self._reader_reconnector_id,
self._id,
self._session_id,
Expand Down Expand Up @@ -572,13 +605,15 @@ def __init__(

self._settings = settings

logger.debug("created ReaderStream id=%s reconnector=%s", self._id, self._reader_reconnector_id)
logger.debug("%s created", self._log_prefix)

@staticmethod
async def create(
reader_reconnector_id: int,
driver: SupportedDriverType,
settings: topic_reader.PublicReaderSettings,
*,
reader_name: str,
) -> "ReaderStream":
stream = GrpcWrapperAsyncIO(StreamReadMessage.FromServer.from_proto)
reader = None
Expand All @@ -590,8 +625,9 @@ async def create(
reader_reconnector_id,
settings,
get_token_function=creds.get_auth_token if creds else None,
reader_name=reader_name,
)
await reader._start(stream, settings._init_message())
await reader._start(stream, settings._init_message(reader_name=reader_name))
except BaseException:
# If create() is interrupted (e.g. reader.close() cancels the connection loop
# mid-reconnect) the in-flight stream is not yet assigned to the reconnector, so
Expand Down Expand Up @@ -623,7 +659,8 @@ async def _start(self, stream: IGrpcWrapperAsyncIO, init_message: StreamReadMess

if isinstance(init_response.server_message, StreamReadMessage.InitResponse):
self._session_id = init_response.server_message.session_id
self._log_prefix = "reader %s stream %s session=%s" % (
self._log_prefix = "topic reader reader_name=%r reader_id=%s stream_id=%s session_id=%s" % (
self._reader_name,
self._reader_reconnector_id,
self._id,
self._session_id,
Expand Down Expand Up @@ -820,7 +857,7 @@ async def _read_messages_loop(self):
"Unexpected message in _read_messages_loop: %s" % type(message.server_message)
)
except issues.UnexpectedGrpcMessage as e:
logger.exception("unexpected message in stream reader: %s" % e)
logger.exception("%s unexpected message in stream reader: %s", self._log_prefix, e)

self._state_changed.set()
except asyncio.CancelledError as e:
Expand Down
Loading
Loading