Kafka error: Erroneous state

Задача состоит в том, чтобы подписаться на топик в кафке и считать данные со смещением по указанной дате. При попытке вызвать метод Seek() получаю ошибку:

Confluent.Kafka.KafkaException: Local: Erroneous state

Вот мой код:

adminClient = new AdminClientBuilder(_kafkaConfig.AsEnumerable()).Build();
var topicMetadata = adminClient.GetMetadata(_config.Topic, TimeSpan.FromSeconds(2));
var partitions = topicMetadata
    .Topics
    .First(x => x.Topic == _config.Topic)
    .Partitions;
var partitionsOffsets = partitions
    .Select(x => new TopicPartitionTimestamp(_config.Topic, x.PartitionId, new Timestamp(_config.OffsetDateUtc)));

consumer = CreateConsumer();

foreach (var p in partitions)
{
    consumer.Assign(new TopicPartition(_config.Topic, p.PartitionId));
}

var offsets = consumer.OffsetsForTimes(partitionsOffsets, TimeSpan.FromSeconds(2));

//await Task.Delay(1000);

foreach (var o in offsets)
{
    consumer.Seek(o);
}

Но если я добавляю ожидание: await Task.Delay(1000); то ошибка не появляется. Как мне правильно задать смещение, чтобы ошибка не появлялась без ожидания?


Ответы (1 шт):

Автор решения: Chloroform

Нашёл верный ответ на вопрос:

var adminClient = new AdminClientBuilder(_kafkaConfig.AsEnumerable()).Build();
var topicMetadata = adminClient.GetMetadata(_config.Topic, TimeSpan.FromSeconds(2));
var partitions = topicMetadata
    .Topics
    .First(x => x.Topic == _config.Topic)
    .Partitions;
var partitionsOffsets = partitions
    .Select(x => new TopicPartitionTimestamp(_config.Topic, x.PartitionId, new Timestamp(_config.OffsetDateUtc)));

consumer = CreateConsumer();

var offsets = consumer.OffsetsForTimes(partitionsOffsets, TimeSpan.FromSeconds(2));

foreach (var o in offsets)
{
    consumer.Assign(o); 
}

Проблема заключается в том, что сначала шла подписка на Partitions, а потом поиск по топику consumer.Seek()

foreach (var p in partitions)
{
    consumer.Assign(new TopicPartition(_config.Topic, p.PartitionId));
}

foreach (var o in offsets)
{
    consumer.Seek(o); 
}

Следующая строка задаёт создаёт необходимые параметры для прдписки на нужные Partitions с нужным смещением:

var offsets = consumer.OffsetsForTimes(partitionsOffsets, TimeSpan.FromSeconds(2));
→ Ссылка