【问题标题】:Can Reactive Kafka Receiver work with non-reactive Elasticsearch client?反应式 Kafka 接收器可以与非反应式 Elasticsearch 客户端一起使用吗?
【发布时间】:2021-04-20 12:59:19
【问题描述】:

下面是一个示例代码,它使用 reactor-kafka 并从一个主题(带有重试逻辑)读取数据,该主题具有通过非响应式生产者发布的记录。在我的doOnNext() 消费者内部,我使用的是非反应性弹性搜索客户端,它为索引中的记录编制索引。所以我有几个问题我仍然不清楚:

  1. 我知道消费者和生产者是独立的解耦系统,但是否建议同时拥有响应式生产者以及其消费者是响应式的?
  2. 如果我使用的是非反应性的东西,在这种情况下是 Elasticsearch 客户端 org.elasticsearch.client.RestClient,代码的 “反应性” 是否有效?如果有或没有,我该如何测试它? (通过“反应性”,我的意思是它的非阻塞 IO 部分,即如果我产生三个反应性消费者并且一个由于某种原因是潜在的,则该线程应该被解除阻塞并用于其他反应性消费者)。
  3. 一般来说,问题是,如果我用响应式客户端包装一些 API,API 是否也应该是响应式的?

public Disposable consumeRecords() {
    long maxAttempts = 3, duration = 10;
    RetryBackoffSpec retrySpec = Retry.backoff(maxAttempts, Duration.ofSeconds(duration)).transientErrors(true);
    Consumer<ReceiverRecord<K, V>> doOnNextConsumer = x -> {
        // use non-reactive elastic search client and index record x
    };

    return KafkaReceiver.create(receiverOptions)
            .receive()
            .doOnNext(record -> {
                try {
                    // calling the non-reactive consumer
                    doOnNextConsumer.accept(record);
                } catch (Exception e) {
                    throw new ReceiverRecordException(record, e);
                }
                record.receiverOffset().acknowledge();
            })
            .doOnError(t -> log.error("Error occurred: ", t))
            .retryWhen(retrySpec)
            .onErrorContinue((e, record) -> {
                ReceiverRecordException receiverRecordException = (ReceiverRecordException) e;
                log.error("Retries exhausted for: " + receiverRecordException);
                receiverRecordException.getRecord().receiverOffset().acknowledge();
            })
            .repeat()
            .subscribe();
}

【问题讨论】:

    标签: elasticsearch apache-kafka reactive-programming producer-consumer reactor-kafka


    【解决方案1】:

    对它有所了解。

    Reactive KafkaReceiver 会在内部调用一些 API;如果那个 API 是 阻塞 API 那么即使 KafkaReceiver 是“反应式”的,非阻塞 IO 也不会工作,并且接收器线程将被阻塞,因为你正在调用 阻塞 API / 非反应式 API

    您可以通过创建一个简单的服务器(它会阻止某个时间/睡眠的调用)并从该接收器调用该服务器来测试这一点

    【讨论】:

      猜你喜欢
      • 2021-07-20
      • 2021-09-12
      • 2021-12-25
      • 2018-04-22
      • 1970-01-01
      • 2019-10-07
      • 2019-06-01
      • 2021-11-27
      • 2022-08-04
      相关资源
      最近更新 更多