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.describe_system_view` (sync and async) returning a dedicated `SystemViewSchemeEntry` with the system view id, name, columns, primary key and attributes — a full description of system view objects that `describe_table` does not provide
* Add `pool_id` parameter to `QuerySession.execute`, `QueryTxContext.execute`, and `QuerySessionPool.execute_with_retries` to route queries to a specific resource pool

## 3.31.1 ##
Expand Down
7 changes: 7 additions & 0 deletions docs/apireference.rst
Original file line number Diff line number Diff line change
Expand Up @@ -372,6 +372,13 @@ TableSchemeEntry
:members:
:undoc-members:

SystemViewSchemeEntry
^^^^^^^^^^^^^^^^^^^^^

.. autoclass:: ydb.SystemViewSchemeEntry
:members:
:undoc-members:

DescribeTableSettings
^^^^^^^^^^^^^^^^^^^^^

Expand Down
22 changes: 22 additions & 0 deletions docs/table.rst
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,24 @@ Pass :class:`~ydb.DescribeTableSettings` to request additional detail:
entry = driver.table_client.describe_table("/local/users", settings)


Describing a System View
^^^^^^^^^^^^^^^^^^^^^^^^

System views (the ``.sys`` schema objects) are described with ``describe_system_view``,
which returns a :class:`~ydb.SystemViewSchemeEntry` with the system view id and name,
its columns, primary key, and attributes. Unlike ``describe_table`` this is a dedicated
description that omits the table-only fields that do not apply to system views:

.. code-block:: python

entry = driver.table_client.describe_system_view("/local/.sys/nodes")

print(entry.sys_view_name, entry.sys_view_id)
for column in entry.columns:
print(column.name, column.type)
print(entry.primary_key)


Altering a Table
^^^^^^^^^^^^^^^^

Expand Down Expand Up @@ -556,6 +574,10 @@ on an async driver. All I/O methods become coroutines:
for col in entry.columns:
print(col.name, col.type)

# Describe a system view
sys_view = await driver.table_client.describe_system_view("/local/.sys/nodes")
print(sys_view.sys_view_name, sys_view.primary_key)

await driver.table_client.drop_table("/local/users")

asyncio.run(main())
Expand Down
16 changes: 16 additions & 0 deletions tests/aio/test_table_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -128,3 +128,19 @@ async def test_rename_index(self, driver: ydb.aio.Driver):
description = await client.describe_table(table_name)
assert len(description.indexes) == 1
assert description.indexes[0].name == "index2"

@pytest.mark.asyncio
async def test_describe_system_view(self, driver: ydb.aio.Driver):
client = driver.table_client

with pytest.raises(ydb.SchemeError):
await client.describe_system_view("/local/.sys/does_not_exist")

entry = await client.describe_system_view("/local/.sys/nodes")

assert isinstance(entry, ydb.SystemViewSchemeEntry)
assert entry.is_sysview()
assert entry.sys_view_name == "nodes"
assert entry.sys_view_id > 0
assert entry.primary_key == ["NodeId"]
assert "NodeId" in [column.name for column in entry.columns]
15 changes: 15 additions & 0 deletions tests/table/test_table_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -168,3 +168,18 @@ def test_rename_index(self, driver_sync: ydb.Driver):
description = client.describe_table(table_name)
assert len(description.indexes) == 1
assert description.indexes[0].name == "index2"

def test_describe_system_view(self, driver_sync: ydb.Driver):
client = driver_sync.table_client

with pytest.raises(ydb.SchemeError):
client.describe_system_view("/local/.sys/does_not_exist")

entry = client.describe_system_view("/local/.sys/nodes")

assert isinstance(entry, ydb.SystemViewSchemeEntry)
assert entry.is_sysview()
assert entry.sys_view_name == "nodes"
assert entry.sys_view_id > 0
assert entry.primary_key == ["NodeId"]
assert "NodeId" in [column.name for column in entry.columns]
1 change: 1 addition & 0 deletions ydb/_apis.py
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,7 @@ class TableService(object):
CopyTables = "CopyTables"
RenameTables = "RenameTables"
DescribeTable = "DescribeTable"
DescribeSystemView = "DescribeSystemView"
CreateSession = "CreateSession"
DeleteSession = "DeleteSession"
ExecuteSchemeQuery = "ExecuteSchemeQuery"
Expand Down
21 changes: 21 additions & 0 deletions ydb/_session_impl.py
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,21 @@ def wrap_describe_table_response(rpc_state, response_pb, sesssion_state, scheme_
)


def wrap_describe_system_view_response(rpc_state, response_pb, scheme_entry_cls):
issues._process_response(response_pb.operation)
message = _apis.ydb_table.DescribeSystemViewResult()
response_pb.operation.result.Unpack(message)
return scheme._wrap_scheme_entry(
message.self,
scheme_entry_cls,
message.sys_view_id,
message.sys_view_name,
message.columns,
message.primary_key,
message.attributes,
)


def explicit_partitions_factory(primary_key, columns, split_points):
column_types = {}
pk = set(primary_key)
Expand Down Expand Up @@ -251,6 +266,12 @@ def describe_table_request_factory(session_state, path, settings=None):
return request


def describe_system_view_request_factory(path, settings=None):
request = _apis.ydb_table.DescribeSystemViewRequest()
request.path = path
return request


def alter_table_request_factory(
session_state,
path,
Expand Down
3 changes: 3 additions & 0 deletions ydb/aio/table.py
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,9 @@ def session(self):
async def bulk_upsert(self, *args, **kwargs): # pylint: disable=W0236
return await super().bulk_upsert(*args, **kwargs)

async def describe_system_view(self, path, settings=None): # pylint: disable=W0236
return await super().describe_system_view(path, settings)

async def scan_query(self, query, parameters=None, settings=None): # pylint: disable=W0236
request = _scan_query_request_factory(query, parameters, settings)
response = await self._driver(
Expand Down
57 changes: 57 additions & 0 deletions ydb/table.py
Original file line number Diff line number Diff line change
Expand Up @@ -1239,6 +1239,25 @@ def bulk_upsert(self, table_path, rows, column_types, settings=None):
(),
)

def describe_system_view(self, path, settings=None):
# type: (str, ydb.BaseRequestSettings) -> Any
"""
Returns a full description of a system view by the provided path.

:param path: A system view path
:param settings: A request settings

:return: SystemViewSchemeEntry describing the system view
"""
Comment thread
vgvoleg marked this conversation as resolved.
return self._driver(
_session_impl.describe_system_view_request_factory(path, settings),
_apis.TableService.Stub,
_apis.TableService.DescribeSystemView,
_session_impl.wrap_describe_system_view_response,
settings,
(SystemViewSchemeEntry,),
)


class TableClient(BaseTableClient["SyncDriver"]):
def __init__(self, driver: "SyncDriver", table_client_settings: Optional[TableClientSettings] = None) -> None:
Expand Down Expand Up @@ -1274,6 +1293,17 @@ def async_bulk_upsert(self, table_path, rows, column_types, settings=None):
(),
)

@_utilities.wrap_async_call_exceptions
def async_describe_system_view(self, path, settings=None):
return self._driver.future(
_session_impl.describe_system_view_request_factory(path, settings),
_apis.TableService.Stub,
_apis.TableService.DescribeSystemView,
_session_impl.wrap_describe_system_view_response,
settings,
(SystemViewSchemeEntry,),
)

def _init_pool_if_needed(self) -> None:
if self._pool is None:
self._pool = SessionPool(self._driver, 10)
Expand Down Expand Up @@ -1650,6 +1680,33 @@ def __init__(
self.attributes = attributes


class SystemViewSchemeEntry(scheme.SchemeEntry):
def __init__(
self,
name,
owner,
type,
effective_permissions,
permissions,
size_bytes,
sys_view_id,
sys_view_name,
columns,
primary_key,
attributes,
*args,
**kwargs
):
super(SystemViewSchemeEntry, self).__init__(
name, owner, type, effective_permissions, permissions, size_bytes, *args, **kwargs
)
self.sys_view_id = sys_view_id
self.sys_view_name = sys_view_name
self.columns = [Column(column.name, convert.type_to_native(column.type), column.family) for column in columns]
self.primary_key = [pk for pk in primary_key]
self.attributes = attributes


class RenameItem:
def __init__(self, source_path, destination_path, replace_destination=False):
self._source_path = source_path
Expand Down
66 changes: 65 additions & 1 deletion ydb/table_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,8 @@
import pytest

from unittest import mock
from . import issues, convert, types, _apis
from . import issues, convert, types, _apis, scheme, _session_impl
from .table import SystemViewSchemeEntry, TableClient

from .retries import (
retry_operation_impl,
Expand Down Expand Up @@ -272,3 +273,66 @@ def check_retriable_error(err_type, backoff):
check_unretriable_error(TestException, False)
with mock.patch.object(retry_once_settings, "idempotent", True):
check_unretriable_error(TestException, False)


def _build_describe_system_view_response():
result = _apis.ydb_table.DescribeSystemViewResult()
result.self.name = "partition_stats"
result.self.type = scheme.SchemeEntryType.SYS_VIEW
result.sys_view_id = 42
result.sys_view_name = "partition_stats"
column = result.columns.add()
column.name = "OwnerId"
column.type.type_id = types.PrimitiveType.Uint64._idn_
result.primary_key.append("OwnerId")
result.attributes["origin"] = "test"

response = _apis.ydb_table.DescribeSystemViewResponse()
response.operation.status = _apis.StatusIds.SUCCESS
response.operation.result.Pack(result)
return response


def test_describe_system_view_request_factory():
request = _session_impl.describe_system_view_request_factory("/local/.sys/partition_stats")
assert request.path == "/local/.sys/partition_stats"
# DescribeSystemView is stateless: the request carries no session id.
assert not hasattr(request, "session_id")


def test_wrap_describe_system_view_response():
entry = _session_impl.wrap_describe_system_view_response(
None, _build_describe_system_view_response(), SystemViewSchemeEntry
)

assert isinstance(entry, SystemViewSchemeEntry)
assert entry.name == "partition_stats"
assert entry.is_sysview()
assert entry.sys_view_id == 42
assert entry.sys_view_name == "partition_stats"
assert entry.primary_key == ["OwnerId"]
assert [column.name for column in entry.columns] == ["OwnerId"]
assert dict(entry.attributes) == {"origin": "test"}


def test_wrap_describe_system_view_response_raises_on_error():
response = _apis.ydb_table.DescribeSystemViewResponse()
response.operation.status = _apis.StatusIds.SCHEME_ERROR
with pytest.raises(issues.SchemeError):
_session_impl.wrap_describe_system_view_response(None, response, SystemViewSchemeEntry)


def test_async_describe_system_view():
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_describe_system_view_response(), *wrap_args)

driver = _FakeSyncDriver()
entry = TableClient(driver).async_describe_system_view("/local/.sys/partition_stats")

assert driver.method == _apis.TableService.DescribeSystemView
assert driver.request.path == "/local/.sys/partition_stats"
assert isinstance(entry, SystemViewSchemeEntry)
assert entry.sys_view_name == "partition_stats"
Loading