【问题标题】:How receive data from kafka from specific date如何从特定日期接收来自 kafka 的数据
【发布时间】:2021-03-02 03:30:29
【问题描述】:

我想订阅主题并接收特定日期的数据。方法调用抛出异常:

Confluent.Kafka.KafkaException:本地:错误状态

我的代码:

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);。方法Seek() 不会抛出异常。如何在没有Task.Delay 的情况下按日期设置偏移量?

【问题讨论】:

    标签: c# apache-kafka kafka-consumer-api


    【解决方案1】:

    下一个代码工作正确!

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

    【讨论】:

    • 你能解释一下吗,你如何阅读基于 dateTime 的 kafka 主题以及这里的一些代码?提前致谢
    猜你喜欢
    • 1970-01-01
    • 2016-03-26
    • 2018-02-08
    • 2011-02-08
    • 1970-01-01
    • 2021-11-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多