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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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

Expand Down
3 changes: 2 additions & 1 deletion docs/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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.


Expand Down
49 changes: 47 additions & 2 deletions docs/table.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
---------------

Expand Down Expand Up @@ -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:
Expand Down
67 changes: 67 additions & 0 deletions tests/aio/test_table_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
66 changes: 66 additions & 0 deletions tests/table/test_table_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
1 change: 1 addition & 0 deletions ydb/_apis.py
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,7 @@ class TableService(object):
KeepAlive = "KeepAlive"
StreamReadTable = "StreamReadTable"
BulkUpsert = "BulkUpsert"
ReadRows = "ReadRows"


class TopicService(object):
Expand Down
14 changes: 14 additions & 0 deletions ydb/_session_impl.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 3 additions & 0 deletions ydb/aio/table.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
48 changes: 48 additions & 0 deletions ydb/table.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
"""
Expand Down Expand Up @@ -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(
Expand Down
Loading
Loading