diff --git a/CHANGELOG.md b/CHANGELOG.md index c04c491e..82d47f40 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,4 @@ +* 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` * Fix query-session gauges after session invalidation and metrics reconfiguration: closed sessions no longer appear as negative `used` or phantom `idle` sessions, provider replacement preserves instrumented pool state without leaking observations to old providers, and pools created with metrics disabled use zero-cost shared no-op lifecycle instrumentation diff --git a/docs/index.rst b/docs/index.rst index 94c7c6d6..66970c20 100644 --- a/docs/index.rst +++ b/docs/index.rst @@ -90,7 +90,8 @@ Table Service The :doc:`table` page covers ``driver.table_client`` — the lower-level API for operations that cannot be expressed in YQL: creating tables with custom partitioning, TTL, secondary indexes, and column families; bulk loading data with ``bulk_upsert``; -and streaming full-table reads with ``read_table`` or ``scan_query``. Use this +point reads by primary key with ``read_rows``; and streaming full-table reads with +``read_table`` or ``scan_query``. Use this alongside the Query service when you need fine-grained schema or data-loading control. diff --git a/docs/table.rst b/docs/table.rst index 6fff47d7..7f92cdc3 100644 --- a/docs/table.rst +++ b/docs/table.rst @@ -3,8 +3,8 @@ Table Service The Table service is the legacy API for schema management and bulk data operations. Use it when you need operations that are not available through YQL — creating tables -with fine-grained partitioning, bulk loading data, streaming full table scans, or -managing secondary indexes programmatically. +with fine-grained partitioning, bulk loading data, point reads by primary key, +streaming full table scans, or managing secondary indexes programmatically. For running queries use :doc:`query` instead. The Table service does not replace the Query service; the two are complementary. @@ -430,6 +430,43 @@ attributes match the column names in ``BulkUpsertColumns``. via the Query service when transactional semantics matter. +Point Reads +----------- + +``read_rows`` reads specific rows by primary key without a transaction or a +session. It is faster than a ``SELECT`` for exact key lookups and returns a +single :class:`~ydb.convert._ResultSet`. + +Describe the primary key columns with :class:`~ydb.BulkUpsertColumns` (or +:class:`~ydb.StructType`) and pass a list of key structures. Missing keys are +omitted from the result rather than raising an error: + +.. code-block:: python + + key_types = ( + ydb.BulkUpsertColumns() + .add_column("id", ydb.PrimitiveType.Uint64) + ) + + result_set = driver.table_client.read_rows( + "/local/users", + [{"id": 1}, {"id": 2}, {"id": 999}], + key_types, + columns=("id", "name"), + ) + + for row in result_set.rows: + print(row.id, row.name) + +Pass ``columns`` to return a subset of fields. Omit it (or pass an empty +iterable) to return every column. + +.. note:: + + ``read_rows`` is not transactional. Prefer a ``SELECT`` via the Query + service when the read must participate in a transaction. + + Streaming Reads --------------- @@ -569,6 +606,14 @@ on an async driver. All I/O methods become coroutines: rows = [...] await driver.table_client.bulk_upsert("/local/users", rows, column_types) + # Point read by primary key + key_types = ydb.BulkUpsertColumns().add_column("id", ydb.PrimitiveType.Uint64) + result_set = await driver.table_client.read_rows( + "/local/users", [{"id": 1}], key_types, columns=("id", "name") + ) + for row in result_set.rows: + print(row.id, row.name) + # Describe entry = await driver.table_client.describe_table("/local/users") for col in entry.columns: diff --git a/tests/aio/test_table_client.py b/tests/aio/test_table_client.py index 860604ce..3347e978 100644 --- a/tests/aio/test_table_client.py +++ b/tests/aio/test_table_client.py @@ -144,3 +144,70 @@ async def test_describe_system_view(self, driver: ydb.aio.Driver): assert entry.sys_view_id > 0 assert entry.primary_key == ["NodeId"] assert "NodeId" in [column.name for column in entry.columns] + + @pytest.mark.asyncio + async def test_read_rows(self, driver: ydb.aio.Driver): + client = driver.table_client + table_name = "/local/testtableclient_read_rows" + try: + await client.drop_table(table_name) + except ydb.SchemeError: + pass + + description = ( + ydb.TableDescription() + .with_primary_keys("key1", "key2") + .with_columns( + ydb.Column("key1", ydb.PrimitiveType.Uint64), + ydb.Column("key2", ydb.PrimitiveType.Uint64), + ydb.Column("value", ydb.OptionalType(ydb.PrimitiveType.Utf8)), + ) + ) + + await client.create_table(table_name, description) + + row_types = ( + ydb.BulkUpsertColumns() + .add_column("key1", ydb.PrimitiveType.Uint64) + .add_column("key2", ydb.PrimitiveType.Uint64) + .add_column("value", ydb.OptionalType(ydb.PrimitiveType.Utf8)) + ) + await client.bulk_upsert( + table_name, + [ + {"key1": 1, "key2": 10, "value": "alice"}, + {"key1": 2, "key2": 20, "value": "bob"}, + {"key1": 3, "key2": 30, "value": "carol"}, + ], + row_types, + ) + + key_types = ( + ydb.BulkUpsertColumns() + .add_column("key1", ydb.PrimitiveType.Uint64) + .add_column("key2", ydb.PrimitiveType.Uint64) + ) + result_set = await client.read_rows( + table_name, + [ + {"key1": 1, "key2": 10}, + {"key1": 3, "key2": 30}, + {"key1": 9, "key2": 90}, + ], + key_types, + ) + + rows = {(row.key1, row.key2, row.value) for row in result_set.rows} + assert rows == {(1, 10, "alice"), (3, 30, "carol")} + + columns_only = await client.read_rows( + table_name, + [{"key1": 2, "key2": 20}], + key_types, + columns=("value",), + ) + assert [column.name for column in columns_only.columns] == ["value"] + assert [row.value for row in columns_only.rows] == ["bob"] + + with pytest.raises(ydb.SchemeError): + await client.read_rows("/local/does_not_exist", [{"key1": 1, "key2": 10}], key_types) diff --git a/tests/table/test_table_client.py b/tests/table/test_table_client.py index 9f9a6fe3..66bf8abc 100644 --- a/tests/table/test_table_client.py +++ b/tests/table/test_table_client.py @@ -183,3 +183,69 @@ def test_describe_system_view(self, driver_sync: ydb.Driver): assert entry.sys_view_id > 0 assert entry.primary_key == ["NodeId"] assert "NodeId" in [column.name for column in entry.columns] + + def test_read_rows(self, driver_sync: ydb.Driver): + client = driver_sync.table_client + table_name = "/local/testtableclient_read_rows" + try: + client.drop_table(table_name) + except ydb.SchemeError: + pass + + description = ( + ydb.TableDescription() + .with_primary_keys("key1", "key2") + .with_columns( + ydb.Column("key1", ydb.PrimitiveType.Uint64), + ydb.Column("key2", ydb.PrimitiveType.Uint64), + ydb.Column("value", ydb.OptionalType(ydb.PrimitiveType.Utf8)), + ) + ) + + client.create_table(table_name, description) + + row_types = ( + ydb.BulkUpsertColumns() + .add_column("key1", ydb.PrimitiveType.Uint64) + .add_column("key2", ydb.PrimitiveType.Uint64) + .add_column("value", ydb.OptionalType(ydb.PrimitiveType.Utf8)) + ) + client.bulk_upsert( + table_name, + [ + {"key1": 1, "key2": 10, "value": "alice"}, + {"key1": 2, "key2": 20, "value": "bob"}, + {"key1": 3, "key2": 30, "value": "carol"}, + ], + row_types, + ) + + key_types = ( + ydb.BulkUpsertColumns() + .add_column("key1", ydb.PrimitiveType.Uint64) + .add_column("key2", ydb.PrimitiveType.Uint64) + ) + result_set = client.read_rows( + table_name, + [ + {"key1": 1, "key2": 10}, + {"key1": 3, "key2": 30}, + {"key1": 9, "key2": 90}, + ], + key_types, + ) + + rows = {(row.key1, row.key2, row.value) for row in result_set.rows} + assert rows == {(1, 10, "alice"), (3, 30, "carol")} + + columns_only = client.read_rows( + table_name, + [{"key1": 2, "key2": 20}], + key_types, + columns=("value",), + ) + assert [column.name for column in columns_only.columns] == ["value"] + assert [row.value for row in columns_only.rows] == ["bob"] + + with pytest.raises(ydb.SchemeError): + client.read_rows("/local/does_not_exist", [{"key1": 1, "key2": 10}], key_types) diff --git a/ydb/_apis.py b/ydb/_apis.py index f4b6dd53..4748106c 100644 --- a/ydb/_apis.py +++ b/ydb/_apis.py @@ -111,6 +111,7 @@ class TableService(object): KeepAlive = "KeepAlive" StreamReadTable = "StreamReadTable" BulkUpsert = "BulkUpsert" + ReadRows = "ReadRows" class TopicService(object): diff --git a/ydb/_session_impl.py b/ydb/_session_impl.py index 8d9ef076..ab34d2a9 100644 --- a/ydb/_session_impl.py +++ b/ydb/_session_impl.py @@ -402,6 +402,20 @@ def bulk_upsert_request_factory(table, rows, column_types): return request +def read_rows_request_factory(table_path, keys, key_types, columns=None): + request = _apis.ydb_table.ReadRowsRequest() + request.path = table_path + request.keys.MergeFrom(convert.to_typed_value_from_native(types.ListType(key_types).proto, keys)) + if columns: + request.columns.extend(list(columns)) + return request + + +def wrap_read_rows_response(rpc_state, response_pb, table_client_settings=None): + issues._process_response(response_pb) + return convert.ResultSet.from_message(response_pb.result_set, table_client_settings) + + def wrap_read_table_response(response): issues._process_response(response) snapshot = response.snapshot if response.HasField("snapshot") else None diff --git a/ydb/aio/table.py b/ydb/aio/table.py index 70db2d87..decd1d37 100644 --- a/ydb/aio/table.py +++ b/ydb/aio/table.py @@ -174,6 +174,9 @@ def session(self): async def bulk_upsert(self, *args, **kwargs): # pylint: disable=W0236 return await super().bulk_upsert(*args, **kwargs) + async def read_rows(self, *args, **kwargs): # pylint: disable=W0236 + return await super().read_rows(*args, **kwargs) + async def describe_system_view(self, path, settings=None): # pylint: disable=W0236 return await super().describe_system_view(path, settings) diff --git a/ydb/table.py b/ydb/table.py index 348caeb0..dad5788f 100644 --- a/ydb/table.py +++ b/ydb/table.py @@ -1195,6 +1195,20 @@ def bulk_upsert(self, table_path, rows, column_types, settings=None): """ pass + @abstractmethod + def read_rows(self, table_path, keys, key_types, columns=None, settings=None): + """ + Read specified keys non-transactionally from a single table. + + :param table_path: A table path. + :param keys: A list of structures matching the primary key. + :param key_types: Primary key column types. + :param columns: Optional iterable of column names to return. + :param settings: Request settings. + + """ + pass + class BaseTableClient(ITableClient, Generic[DriverT]): _driver: DriverT @@ -1240,6 +1254,28 @@ def bulk_upsert(self, table_path, rows, column_types, settings=None): (), ) + def read_rows(self, table_path, keys, key_types, columns=None, settings=None): + # type: (str, list, ydb.AbstractTypeBuilder, typing.Optional[list], ydb.BaseRequestSettings) -> Any + """ + Read specified keys non-transactionally from a single table. + + :param table_path: A table path. + :param keys: A list of structures matching the primary key. + :param key_types: Primary key column types. + :param columns: Optional iterable of column names to return. Empty or omitted returns all columns. + :param settings: Request settings. + + :return: ResultSet with matching rows. + """ + return self._driver( + _session_impl.read_rows_request_factory(table_path, keys, key_types, columns), + _apis.TableService.Stub, + _apis.TableService.ReadRows, + _session_impl.wrap_read_rows_response, + settings, + (self._table_client_settings,), + ) + def describe_system_view(self, path, settings=None): # type: (str, ydb.BaseRequestSettings) -> Any """ @@ -1294,6 +1330,18 @@ def async_bulk_upsert(self, table_path, rows, column_types, settings=None): (), ) + @_utilities.wrap_async_call_exceptions + def async_read_rows(self, table_path, keys, key_types, columns=None, settings=None): + # type: (str, list, ydb.AbstractTypeBuilder, typing.Optional[list], ydb.BaseRequestSettings) -> Any + return self._driver.future( + _session_impl.read_rows_request_factory(table_path, keys, key_types, columns), + _apis.TableService.Stub, + _apis.TableService.ReadRows, + _session_impl.wrap_read_rows_response, + settings, + (self._table_client_settings,), + ) + @_utilities.wrap_async_call_exceptions def async_describe_system_view(self, path, settings=None): return self._driver.future( diff --git a/ydb/table_test.py b/ydb/table_test.py index 23a71413..e6a33af6 100644 --- a/ydb/table_test.py +++ b/ydb/table_test.py @@ -338,6 +338,98 @@ def future(self, request, stub, method, wrap_fn, settings, wrap_args, *rest): assert entry.sys_view_name == "partition_stats" +def _read_rows_key_types(): + return types.BulkUpsertColumns().add_column("id", types.PrimitiveType.Uint64) + + +def _build_read_rows_response(): + response = _apis.ydb_table.ReadRowsResponse() + response.status = _apis.StatusIds.SUCCESS + id_column = response.result_set.columns.add() + id_column.name = "id" + id_column.type.type_id = types.PrimitiveType.Uint64._idn_ + value_column = response.result_set.columns.add() + value_column.name = "value" + value_column.type.type_id = types.PrimitiveType.Utf8._idn_ + row = response.result_set.rows.add() + row.items.add().uint64_value = 1 + row.items.add().text_value = "alice" + return response + + +def test_read_rows_request_factory(): + keys = [{"id": 1}, {"id": 2}] + request = _session_impl.read_rows_request_factory( + "/local/users", keys, _read_rows_key_types(), columns=("id", "value") + ) + + assert request.path == "/local/users" + # ReadRows is stateless: the request carries no session id. + assert request.session_id == "" + assert list(request.columns) == ["id", "value"] + assert request.keys.type.list_type.item.struct_type.members[0].name == "id" + assert [item.items[0].uint64_value for item in request.keys.value.items] == [1, 2] + + +def test_read_rows_request_factory_without_columns(): + request = _session_impl.read_rows_request_factory("/local/users", [{"id": 1}], _read_rows_key_types()) + assert list(request.columns) == [] + + +def test_wrap_read_rows_response(): + result_set = _session_impl.wrap_read_rows_response(None, _build_read_rows_response()) + + assert [column.name for column in result_set.columns] == ["id", "value"] + assert len(result_set.rows) == 1 + assert result_set.rows[0].id == 1 + assert result_set.rows[0].value == "alice" + + +def test_wrap_read_rows_response_raises_on_error(): + response = _apis.ydb_table.ReadRowsResponse() + response.status = _apis.StatusIds.SCHEME_ERROR + with pytest.raises(issues.SchemeError): + _session_impl.wrap_read_rows_response(None, response) + + +def test_async_read_rows(): + class _FakeSyncDriver: + def future(self, request, stub, method, wrap_fn, settings, wrap_args, *rest): + self.request = request + self.method = method + return wrap_fn(None, _build_read_rows_response(), *wrap_args) + + driver = _FakeSyncDriver() + result_set = TableClient(driver).async_read_rows( + "/local/users", [{"id": 1}], _read_rows_key_types(), columns=("id", "value") + ) + + assert driver.method == _apis.TableService.ReadRows + assert driver.request.path == "/local/users" + assert list(driver.request.columns) == ["id", "value"] + assert result_set.rows[0].id == 1 + assert result_set.rows[0].value == "alice" + + +def test_read_rows(): + class _FakeSyncDriver: + def __call__(self, request, stub, method, wrap_fn, settings, wrap_args, *rest): + self.request = request + self.method = method + return wrap_fn(None, _build_read_rows_response(), *wrap_args) + + driver = _FakeSyncDriver() + result_set = TableClient(driver).read_rows( + "/local/users", [{"id": 1}], _read_rows_key_types(), columns=("id", "value") + ) + + assert driver.method == _apis.TableService.ReadRows + assert driver.request.path == "/local/users" + assert list(driver.request.columns) == ["id", "value"] + assert result_set.rows[0].id == 1 + assert result_set.rows[0].value == "alice" + + def _read_table_session_state(): state = _session_impl.SessionState(TableClientSettings()) state.set_id("test-session-id")