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 `TopicClient.reset_offset` (sync and async) to rewind a consumer's committed offsets on all topic partitions, including inactive ones, to earliest, latest, or the first message written at or after a given time
* 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
7 changes: 7 additions & 0 deletions docs/apireference.rst
Original file line number Diff line number Diff line change
Expand Up @@ -484,6 +484,13 @@ TopicCodec
:members:
:undoc-members:

TopicResetOffset
^^^^^^^^^^^^^^^^

.. autoclass:: ydb.TopicResetOffset
:members:
:undoc-members:

TopicWriterMessage
^^^^^^^^^^^^^^^^^^

Expand Down
28 changes: 28 additions & 0 deletions docs/topic.rst
Original file line number Diff line number Diff line change
Expand Up @@ -507,6 +507,34 @@ from the last committed offset when the reader reconnects.
automatically when the reader is closed (``flush=True`` by default).


Resetting Offsets
^^^^^^^^^^^^^^^^^

``reset_offset`` rewinds a consumer's committed offsets on every topic partition, including
inactive partitions left after a split or merge. Partitions are updated independently: the
call is not atomic, and a failure may still leave some partitions already rewritten. Any
active read session of this consumer is dropped.

.. code-block:: python

import datetime

# Rewind to the beginning or to the end of every partition.
driver.topic_client.reset_offset(topic_path, consumer, to=ydb.TopicResetOffset.EARLIEST)
driver.topic_client.reset_offset(topic_path, consumer, to=ydb.TopicResetOffset.LATEST)

# Rewind to the first message written at or after the given time.
# If there is no such message, the partition end offset is used.
driver.topic_client.reset_offset(
topic_path,
consumer,
to=datetime.datetime(2026, 1, 1, tzinfo=datetime.timezone.utc),
)

# Asynchronous client has the same method.
await driver.topic_client.reset_offset(topic_path, consumer, to=ydb.TopicResetOffset.EARLIEST)


Handling Partition Rebalancing
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^

Expand Down
5 changes: 5 additions & 0 deletions examples/topic/reader_example.py
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,11 @@ def process_batch(batch):
process_batch(batch)


def reset_consumer_offset(db: ydb.Driver):
db.topic_client.reset_offset("/local/topic", "consumer", to=ydb.TopicResetOffset.EARLIEST)
db.topic_client.reset_offset("/local/topic", "consumer", to=ydb.TopicResetOffset.LATEST)


def handle_partition_graceful_stop_batch(reader: ydb.TopicReader):
# no special handle, but batch will contain less than prefer count messages
while True:
Expand Down
153 changes: 153 additions & 0 deletions tests/topics/test_control_plane.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import datetime
import os.path

import pytest
Expand All @@ -6,6 +7,14 @@
from ydb import issues


def _committed_offset(description, partition_id=0):
for partition in description.partitions:
if partition.partition_id == partition_id:
assert partition.partition_consumer_stats is not None
return partition.partition_consumer_stats.committed_offset
raise AssertionError("partition %s not found" % partition_id)


@pytest.mark.asyncio
class TestTopicClientControlPlaneAsyncIO:
async def test_create_topic(self, driver, database):
Expand Down Expand Up @@ -148,6 +157,78 @@ async def test_alter_auto_partitioning_settings(self, driver, topic_path):

assert topic_after.auto_partitioning_settings == expected

async def test_reset_offset(self, driver, topic_with_messages, topic_consumer):
client = driver.topic_client
path = topic_with_messages

await client.reset_offset(path, topic_consumer, to=ydb.TopicResetOffset.LATEST)
description = await client.describe_consumer(path, topic_consumer, include_stats=True)
assert _committed_offset(description) == 4

await client.reset_offset(path, topic_consumer, to=ydb.TopicResetOffset.LATEST)
description = await client.describe_consumer(path, topic_consumer, include_stats=True)
assert _committed_offset(description) == 4

await client.reset_offset(path, topic_consumer, to=ydb.TopicResetOffset.EARLIEST)
description = await client.describe_consumer(path, topic_consumer, include_stats=True)
assert _committed_offset(description) == 0

async with client.reader(path, topic_consumer) as reader:
message = await reader.receive_message()
assert message.data.decode() == "123"

written_at = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(hours=1)
await client.reset_offset(path, topic_consumer, to=written_at)
description = await client.describe_consumer(path, topic_consumer, include_stats=True)
assert _committed_offset(description) == 4

async def test_reset_offset_missing_topic_and_consumer(self, driver, topic_path, topic_consumer):
client = driver.topic_client

with pytest.raises(issues.SchemeError):
await client.reset_offset(topic_path + "-missing", topic_consumer, to=ydb.TopicResetOffset.EARLIEST)

with pytest.raises(issues.SchemeError):
await client.reset_offset(topic_path, "no-such-consumer", to=ydb.TopicResetOffset.EARLIEST)

async def test_reset_offset_leaves_other_consumer(self, driver, database):
client = driver.topic_client
path = database + "/reset-offset-two-consumers"
try:
await client.drop_topic(path)
except issues.SchemeError:
pass

await client.create_topic(path, consumers=["consumer-a", "consumer-b"])
async with client.writer(path) as writer:
await writer.write_with_ack(
[
ydb.TopicWriterMessage(data=b"m1"),
ydb.TopicWriterMessage(data=b"m2"),
]
)

await client.reset_offset(path, "consumer-a", to=ydb.TopicResetOffset.LATEST)

consumer_a = await client.describe_consumer(path, "consumer-a", include_stats=True)
consumer_b = await client.describe_consumer(path, "consumer-b", include_stats=True)
assert _committed_offset(consumer_a) == 2
assert _committed_offset(consumer_b) == 0

async def test_reset_offset_all_partitions(self, driver, topic_with_two_partitions_path, topic_consumer):
client = driver.topic_client
path = topic_with_two_partitions_path

for partition_id in (0, 1):
async with client.writer(path, partition_id=partition_id) as writer:
await writer.write_with_ack(ydb.TopicWriterMessage(data=b"m"))

await client.reset_offset(path, topic_consumer, to=ydb.TopicResetOffset.LATEST)

description = await client.describe_consumer(path, topic_consumer, include_stats=True)
assert _committed_offset(description, 0) == 1
assert _committed_offset(description, 1) == 1


class TestTopicClientControlPlane:
def test_create_topic(self, driver_sync, database):
Expand Down Expand Up @@ -253,3 +334,75 @@ def test_alter_existed_topic(self, driver_sync, topic_path):

topic_after = client.describe_topic(topic_path)
assert topic_after.min_active_partitions == target_min_active_partitions

def test_reset_offset(self, driver_sync, topic_with_messages, topic_consumer):
client = driver_sync.topic_client
path = topic_with_messages

client.reset_offset(path, topic_consumer, to=ydb.TopicResetOffset.LATEST)
description = client.describe_consumer(path, topic_consumer, include_stats=True)
assert _committed_offset(description) == 4

client.reset_offset(path, topic_consumer, to=ydb.TopicResetOffset.LATEST)
description = client.describe_consumer(path, topic_consumer, include_stats=True)
assert _committed_offset(description) == 4

client.reset_offset(path, topic_consumer, to=ydb.TopicResetOffset.EARLIEST)
description = client.describe_consumer(path, topic_consumer, include_stats=True)
assert _committed_offset(description) == 0

with client.reader(path, topic_consumer) as reader:
message = reader.receive_message()
assert message.data.decode() == "123"

written_at = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(hours=1)
client.reset_offset(path, topic_consumer, to=written_at)
description = client.describe_consumer(path, topic_consumer, include_stats=True)
assert _committed_offset(description) == 4

def test_reset_offset_missing_topic_and_consumer(self, driver_sync, topic_path, topic_consumer):
client = driver_sync.topic_client

with pytest.raises(issues.SchemeError):
client.reset_offset(topic_path + "-missing", topic_consumer, to=ydb.TopicResetOffset.EARLIEST)

with pytest.raises(issues.SchemeError):
client.reset_offset(topic_path, "no-such-consumer", to=ydb.TopicResetOffset.EARLIEST)

def test_reset_offset_leaves_other_consumer(self, driver_sync, database):
client = driver_sync.topic_client
path = database + "/reset-offset-two-consumers-sync"
try:
client.drop_topic(path)
except issues.SchemeError:
pass

client.create_topic(path, consumers=["consumer-a", "consumer-b"])
with client.writer(path) as writer:
writer.write_with_ack(
[
ydb.TopicWriterMessage(data=b"m1"),
ydb.TopicWriterMessage(data=b"m2"),
]
)

client.reset_offset(path, "consumer-a", to=ydb.TopicResetOffset.LATEST)

consumer_a = client.describe_consumer(path, "consumer-a", include_stats=True)
consumer_b = client.describe_consumer(path, "consumer-b", include_stats=True)
assert _committed_offset(consumer_a) == 2
assert _committed_offset(consumer_b) == 0

def test_reset_offset_all_partitions(self, driver_sync, topic_with_two_partitions_path, topic_consumer):
client = driver_sync.topic_client
path = topic_with_two_partitions_path

for partition_id in (0, 1):
with client.writer(path, partition_id=partition_id) as writer:
writer.write_with_ack(ydb.TopicWriterMessage(data=b"m"))

client.reset_offset(path, topic_consumer, to=ydb.TopicResetOffset.LATEST)

description = client.describe_consumer(path, topic_consumer, include_stats=True)
assert _committed_offset(description, 0) == 1
assert _committed_offset(description, 1) == 1
1 change: 1 addition & 0 deletions ydb/_apis.py
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,7 @@ class TopicService(object):
StreamWrite = "StreamWrite"
UpdateOffsetsInTransaction = "UpdateOffsetsInTransaction"
CommitOffset = "CommitOffset"
ResetOffset = "ResetOffset"


class QueryService(object):
Expand Down
29 changes: 29 additions & 0 deletions ydb/_grpc/grpcwrapper/ydb_topic.py
Original file line number Diff line number Diff line change
Expand Up @@ -162,6 +162,35 @@ def to_proto(self) -> ydb_topic_pb2.CommitOffsetRequest:
)


@dataclass
class ResetOffsetRequest(IToProto):
path: str
consumer: str
to: Union[ydb_topic_public_types.PublicResetOffset, datetime.datetime]

def __post_init__(self):
if isinstance(self.to, datetime.datetime):
return
if isinstance(self.to, ydb_topic_public_types.PublicResetOffset):
return
raise TypeError(
"reset offset target must be TopicResetOffset.EARLIEST, TopicResetOffset.LATEST, or datetime.datetime"
)

def to_proto(self) -> ydb_topic_pb2.ResetOffsetRequest:
res = ydb_topic_pb2.ResetOffsetRequest(
path=self.path,
consumer=self.consumer,
)
if isinstance(self.to, datetime.datetime):
res.from_written_at.written_at.FromDatetime(self.to)
elif self.to == ydb_topic_public_types.PublicResetOffset.EARLIEST:
res.earliest.SetInParent()
elif self.to == ydb_topic_public_types.PublicResetOffset.LATEST:
res.latest.SetInParent()
return res


########################################################################################################################
# StreamWrite
########################################################################################################################
Expand Down
14 changes: 13 additions & 1 deletion ydb/_grpc/grpcwrapper/ydb_topic_public_types.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import datetime
import typing
from dataclasses import dataclass, field
from enum import IntEnum
from enum import Enum, IntEnum
from typing import Optional, List, Union, Dict

# Workaround for good IDE and universal for runtime
Expand Down Expand Up @@ -66,6 +66,18 @@ class PublicCodec(int):
ZSTD = 4 # Has not supported codec in standard library


class PublicResetOffset(Enum):
"""
Where to rewind a consumer's committed offsets.

Pass a ``datetime`` to :meth:`TopicClient.reset_offset` instead of this enum
to rewind to the first message with write timestamp greater than or equal to that time.
"""

EARLIEST = "earliest"
LATEST = "latest"


class PublicMeteringMode(IntEnum):
UNSPECIFIED = 0
RESERVED_CAPACITY = 1
Expand Down
25 changes: 24 additions & 1 deletion ydb/_grpc/grpcwrapper/ydb_topic_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,10 @@

from google.protobuf.json_format import MessageToDict

from ydb._grpc.grpcwrapper.ydb_topic import OffsetsRange
from ydb._grpc.grpcwrapper.ydb_topic import OffsetsRange, ResetOffsetRequest
from .ydb_topic import AlterTopicRequest
from .ydb_topic_public_types import (
PublicResetOffset,
AlterTopicRequestParams,
PublicAlterConsumer,
PublicAlterAutoPartitioningSettings,
Expand Down Expand Up @@ -96,3 +97,25 @@ def test_alter_topic_request_from_public_to_proto():
}

assert msg_dict == expected_dict


def test_reset_offset_request_to_proto():
earliest = ResetOffsetRequest(path="topic", consumer="consumer", to=PublicResetOffset.EARLIEST).to_proto()
assert earliest.path == "topic"
assert earliest.consumer == "consumer"
assert earliest.WhichOneof("position") == "earliest"

latest = ResetOffsetRequest(path="topic", consumer="consumer", to=PublicResetOffset.LATEST).to_proto()
assert latest.WhichOneof("position") == "latest"

written_at = datetime.datetime(2026, 1, 2, 3, 4, 5, tzinfo=datetime.timezone.utc)
stamped = ResetOffsetRequest(path="topic", consumer="consumer", to=written_at).to_proto()
assert stamped.WhichOneof("position") == "from_written_at"
assert stamped.from_written_at.written_at.ToSeconds() == int(written_at.timestamp())

try:
ResetOffsetRequest(path="topic", consumer="consumer", to="earliest")
except TypeError:
pass
else:
raise AssertionError("expected TypeError")
448 changes: 361 additions & 87 deletions ydb/_grpc/v3/protos/ydb_topic_pb2.py

Large diffs are not rendered by default.

Loading
Loading