org.apache.kafka.clients.consumer.Consumer.wakeup()方法的使用及代码示例

x33g5p2x  于2022-01-18 转载在 其他  
字(6.5k)|赞(0)|评价(0)|浏览(284)

本文整理了Java中org.apache.kafka.clients.consumer.Consumer.wakeup()方法的一些代码示例,展示了Consumer.wakeup()的具体用法。这些代码示例主要来源于Github/Stackoverflow/Maven等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。Consumer.wakeup()方法的具体详情如下:
包路径:org.apache.kafka.clients.consumer.Consumer
类名称:Consumer
方法名:wakeup

Consumer.wakeup介绍

暂无

代码示例

代码示例来源:origin: apache/nifi

/**
 * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
 */
public void wakeup() {
  kafkaConsumer.wakeup();
}

代码示例来源:origin: apache/nifi

/**
 * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
 */
public void wakeup() {
  kafkaConsumer.wakeup();
}

代码示例来源:origin: apache/nifi

/**
 * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
 */
public void wakeup() {
  kafkaConsumer.wakeup();
}

代码示例来源:origin: apache/nifi

/**
 * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
 */
public void wakeup() {
  kafkaConsumer.wakeup();
}

代码示例来源:origin: apache/nifi

/**
 * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
 */
public void wakeup() {
  kafkaConsumer.wakeup();
}

代码示例来源:origin: openzipkin/brave

@Override public void wakeup() {
 delegate.wakeup();
}

代码示例来源:origin: apache/incubator-gobblin

@Override
public void close()
  throws IOException {
 _close.set(true);
 _consumer.wakeup();
 synchronized (_consumer) {
  closer.close();
 }
}

代码示例来源:origin: confluentinc/ksql

public void close() {
  commandConsumer.wakeup();
  commandConsumer.close();
  commandProducer.close();
 }
}

代码示例来源:origin: confluentinc/ksql

@Test
public void shouldCloseAllResources() {
 // When:
 commandTopic.close();
 //Then:
 final InOrder ordered = inOrder(commandConsumer);
 ordered.verify(commandConsumer).wakeup();
 ordered.verify(commandConsumer).close();
 verify(commandProducer).close();
}

代码示例来源:origin: spring-projects/spring-kafka

@Override
protected void doStop(final Runnable callback) {
  if (isRunning()) {
    this.listenerConsumerFuture.addCallback(new StopCallback(callback));
    setRunning(false);
    this.listenerConsumer.consumer.wakeup();
  }
}

代码示例来源:origin: spring-projects/spring-kafka

@SuppressWarnings("unchecked")
@Test
public void stopContainerAfterException() throws Exception {
  assertThat(this.config.deliveryLatch.await(10, TimeUnit.SECONDS)).isTrue();
  assertThat(this.config.pollLatch.await(10, TimeUnit.SECONDS)).isTrue();
  assertThat(this.config.errorLatch.await(10, TimeUnit.SECONDS)).isTrue();
  assertThat(this.config.closeLatch.await(10, TimeUnit.SECONDS)).isTrue();
  MessageListenerContainer container = this.registry.getListenerContainer(CONTAINER_ID);
  assertThat(container.isRunning()).isFalse();
  InOrder inOrder = inOrder(this.consumer);
  inOrder.verify(this.consumer).subscribe(any(Collection.class), any(ConsumerRebalanceListener.class));
  inOrder.verify(this.consumer).poll(Duration.ofMillis(ContainerProperties.DEFAULT_POLL_TIMEOUT));
  inOrder.verify(this.consumer).wakeup();
  inOrder.verify(this.consumer).unsubscribe();
  inOrder.verify(this.consumer).close();
  inOrder.verifyNoMoreInteractions();
}

代码示例来源:origin: spring-projects/spring-kafka

@SuppressWarnings("unchecked")
@Test
public void stopContainerAfterException() throws Exception {
  assertThat(this.config.deliveryLatch.await(10, TimeUnit.SECONDS)).isTrue();
  assertThat(this.config.pollLatch.await(10, TimeUnit.SECONDS)).isTrue();
  assertThat(this.config.errorLatch.await(10, TimeUnit.SECONDS)).isTrue();
  assertThat(this.config.closeLatch.await(10, TimeUnit.SECONDS)).isTrue();
  MessageListenerContainer container = this.registry.getListenerContainer(CONTAINER_ID);
  assertThat(container.isRunning()).isFalse();
  InOrder inOrder = inOrder(this.consumer);
  inOrder.verify(this.consumer).subscribe(any(Collection.class), any(ConsumerRebalanceListener.class));
  inOrder.verify(this.consumer).poll(Duration.ofMillis(ContainerProperties.DEFAULT_POLL_TIMEOUT));
  inOrder.verify(this.consumer).wakeup();
  inOrder.verify(this.consumer).unsubscribe();
  inOrder.verify(this.consumer).close();
  inOrder.verifyNoMoreInteractions();
  assertThat(this.config.count).isEqualTo(4);
  assertThat(this.config.contents.toArray()).isEqualTo(new String[]
      { "foo", "bar", "baz", "qux" });
}

代码示例来源:origin: spring-projects/spring-kafka

@SuppressWarnings("unchecked")
@Test
public void stopContainerAfterException() throws Exception {
  assertThat(this.config.deliveryLatch.await(10, TimeUnit.SECONDS)).isTrue();
  assertThat(this.config.commitLatch.await(10, TimeUnit.SECONDS)).isTrue();
  assertThat(this.config.pollLatch.await(10, TimeUnit.SECONDS)).isTrue();
  assertThat(this.config.errorLatch.await(10, TimeUnit.SECONDS)).isTrue();
  assertThat(this.config.closeLatch.await(10, TimeUnit.SECONDS)).isTrue();
  MessageListenerContainer container = this.registry.getListenerContainer(CONTAINER_ID);
  assertThat(container.isRunning()).isFalse();
  InOrder inOrder = inOrder(this.consumer);
  inOrder.verify(this.consumer).subscribe(any(Collection.class), any(ConsumerRebalanceListener.class));
  inOrder.verify(this.consumer).poll(Duration.ofMillis(ContainerProperties.DEFAULT_POLL_TIMEOUT));
  inOrder.verify(this.consumer).commitSync(
      Collections.singletonMap(new TopicPartition("foo", 0), new OffsetAndMetadata(1L)));
  inOrder.verify(this.consumer).commitSync(
      Collections.singletonMap(new TopicPartition("foo", 0), new OffsetAndMetadata(2L)));
  inOrder.verify(this.consumer).commitSync(
      Collections.singletonMap(new TopicPartition("foo", 1), new OffsetAndMetadata(1L)));
  inOrder.verify(this.consumer).wakeup();
  inOrder.verify(this.consumer).unsubscribe();
  inOrder.verify(this.consumer).close();
  inOrder.verifyNoMoreInteractions();
  assertThat(this.config.count).isEqualTo(4);
  assertThat(this.config.contents.toArray()).isEqualTo(new String[]
      { "foo", "bar", "baz", "qux" });
}

代码示例来源:origin: org.apache.nifi/nifi-kafka-1-0-processors

/**
 * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
 */
public void wakeup() {
  kafkaConsumer.wakeup();
}

代码示例来源:origin: rayokota/kafka-graphs

@Override
  public void wakeup() {
    kafkaConsumer.wakeup();
  }
}

代码示例来源:origin: org.apache.nifi/nifi-kafka-0-9-processors

/**
 * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
 */
public void wakeup() {
  kafkaConsumer.wakeup();
}

代码示例来源:origin: io.opentracing.contrib/opentracing-kafka-client

@Override
 public void wakeup() {
  consumer.wakeup();
 }
}

代码示例来源:origin: org.apache.nifi/nifi-kafka-0-10-processors

/**
 * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
 */
public void wakeup() {
  kafkaConsumer.wakeup();
}

代码示例来源:origin: spring-projects/spring-kafka

deadLatch.countDown();
  return null;
}).given(consumer).wakeup();
TopicPartitionInitialOffset[] topicPartition = new TopicPartitionInitialOffset[] {
    new TopicPartitionInitialOffset("foo", 0) };

代码示例来源:origin: pentaho/big-data-plugin

public void shutdown() {
 closed.set( true );
 consumer.wakeup();
}

相关文章