【问题标题】:kafka asynchronous send not really asynchronous?kafka异步发送不是真的异步吗?
【发布时间】:2019-11-30 02:44:36
【问题描述】:

我正在使用来自 kafka-client 1.0.0 库的 KafkaProducer,根据文档,Future<RecordMetadata> send(ProducerRecord<K, V> record) 方法将立即返回,但实际上,但看起来没有。这个方法还调用了同一个类中的另一个方法doSend(sn-p见下文),在这个方法内部,它正在等待主题的元数据,我认为这是必要的,因为它与分区等。

/**
 * Implementation of asynchronously send a record to a topic.
 */
private Future<RecordMetadata> doSend(ProducerRecord<K, V> record, Callback callback) {
    TopicPartition tp = null;
    try {
        // first make sure the metadata for the topic is available
        ClusterAndWaitTime clusterAndWaitTime = waitOnMetadata(record.topic(), record.partition(), maxBlockTimeMs);
        long remainingWaitMs = Math.max(0, maxBlockTimeMs - clusterAndWaitTime.waitedOnMetadataMs);
        Cluster cluster = clusterAndWaitTime.cluster;

还有其他完全异步的选项吗?我希望它完全异步的问题是因为如果bootstrap.servers 中的某些服务器没有响应,它将等待基于max.block.ms 的时间,但我实际上并不希望它等待,但相反,我只是希望它返回。

我看到它会立即返回的文档: KafkaProducer java doc

发送是异步的,该方法会立即返回一次 记录已存储在等待记录的缓冲区中 发送。这允许并行发送许多记录而不会阻塞 等待每个之后的响应。

【问题讨论】:

    标签: apache-kafka kafka-producer-api


    【解决方案1】:

    您的分析是正确的 - kafka 有一个(有时)阻塞的“非阻塞”API。 这之前已经提出过 - https://cwiki.apache.org/confluence/display/KAFKA/KIP-286%3A+producer.send%28%29+should+not+block+on+metadata+update - 但从未优先考虑。

    【讨论】:

      【解决方案2】:

      它尽可能地异步。 Kafka 维护一个元数据缓存,该缓存偶尔会更新以使其保持最新状态,在您的场景中,您只需等待该缓存过时或未初始化。缓存初始化后,无需等待。

      如果您的代码有一个必须尽快执行的即将到来的 send(),您可以尝试向生产者发送一个准备好的 partitionsFor() 方法调用,以查看是否无法在需要时强制更新缓存。

      除此之外,总会有可能偶尔等待元数据缓存被刷新。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2011-08-05
        • 2014-04-25
        • 2014-09-19
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2018-02-11
        相关资源
        最近更新 更多