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));