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

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

本文整理了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

  1. /**
  2. * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
  3. */
  4. public void wakeup() {
  5. kafkaConsumer.wakeup();
  6. }

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

  1. /**
  2. * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
  3. */
  4. public void wakeup() {
  5. kafkaConsumer.wakeup();
  6. }

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

  1. /**
  2. * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
  3. */
  4. public void wakeup() {
  5. kafkaConsumer.wakeup();
  6. }

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

  1. /**
  2. * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
  3. */
  4. public void wakeup() {
  5. kafkaConsumer.wakeup();
  6. }

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

  1. /**
  2. * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
  3. */
  4. public void wakeup() {
  5. kafkaConsumer.wakeup();
  6. }

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

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

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

  1. @Override
  2. public void close()
  3. throws IOException {
  4. _close.set(true);
  5. _consumer.wakeup();
  6. synchronized (_consumer) {
  7. closer.close();
  8. }
  9. }

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

  1. public void close() {
  2. commandConsumer.wakeup();
  3. commandConsumer.close();
  4. commandProducer.close();
  5. }
  6. }

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

  1. @Test
  2. public void shouldCloseAllResources() {
  3. // When:
  4. commandTopic.close();
  5. //Then:
  6. final InOrder ordered = inOrder(commandConsumer);
  7. ordered.verify(commandConsumer).wakeup();
  8. ordered.verify(commandConsumer).close();
  9. verify(commandProducer).close();
  10. }

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

  1. @Override
  2. protected void doStop(final Runnable callback) {
  3. if (isRunning()) {
  4. this.listenerConsumerFuture.addCallback(new StopCallback(callback));
  5. setRunning(false);
  6. this.listenerConsumer.consumer.wakeup();
  7. }
  8. }

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

  1. @SuppressWarnings("unchecked")
  2. @Test
  3. public void stopContainerAfterException() throws Exception {
  4. assertThat(this.config.deliveryLatch.await(10, TimeUnit.SECONDS)).isTrue();
  5. assertThat(this.config.pollLatch.await(10, TimeUnit.SECONDS)).isTrue();
  6. assertThat(this.config.errorLatch.await(10, TimeUnit.SECONDS)).isTrue();
  7. assertThat(this.config.closeLatch.await(10, TimeUnit.SECONDS)).isTrue();
  8. MessageListenerContainer container = this.registry.getListenerContainer(CONTAINER_ID);
  9. assertThat(container.isRunning()).isFalse();
  10. InOrder inOrder = inOrder(this.consumer);
  11. inOrder.verify(this.consumer).subscribe(any(Collection.class), any(ConsumerRebalanceListener.class));
  12. inOrder.verify(this.consumer).poll(Duration.ofMillis(ContainerProperties.DEFAULT_POLL_TIMEOUT));
  13. inOrder.verify(this.consumer).wakeup();
  14. inOrder.verify(this.consumer).unsubscribe();
  15. inOrder.verify(this.consumer).close();
  16. inOrder.verifyNoMoreInteractions();
  17. }

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

  1. @SuppressWarnings("unchecked")
  2. @Test
  3. public void stopContainerAfterException() throws Exception {
  4. assertThat(this.config.deliveryLatch.await(10, TimeUnit.SECONDS)).isTrue();
  5. assertThat(this.config.pollLatch.await(10, TimeUnit.SECONDS)).isTrue();
  6. assertThat(this.config.errorLatch.await(10, TimeUnit.SECONDS)).isTrue();
  7. assertThat(this.config.closeLatch.await(10, TimeUnit.SECONDS)).isTrue();
  8. MessageListenerContainer container = this.registry.getListenerContainer(CONTAINER_ID);
  9. assertThat(container.isRunning()).isFalse();
  10. InOrder inOrder = inOrder(this.consumer);
  11. inOrder.verify(this.consumer).subscribe(any(Collection.class), any(ConsumerRebalanceListener.class));
  12. inOrder.verify(this.consumer).poll(Duration.ofMillis(ContainerProperties.DEFAULT_POLL_TIMEOUT));
  13. inOrder.verify(this.consumer).wakeup();
  14. inOrder.verify(this.consumer).unsubscribe();
  15. inOrder.verify(this.consumer).close();
  16. inOrder.verifyNoMoreInteractions();
  17. assertThat(this.config.count).isEqualTo(4);
  18. assertThat(this.config.contents.toArray()).isEqualTo(new String[]
  19. { "foo", "bar", "baz", "qux" });
  20. }

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

  1. @SuppressWarnings("unchecked")
  2. @Test
  3. public void stopContainerAfterException() throws Exception {
  4. assertThat(this.config.deliveryLatch.await(10, TimeUnit.SECONDS)).isTrue();
  5. assertThat(this.config.commitLatch.await(10, TimeUnit.SECONDS)).isTrue();
  6. assertThat(this.config.pollLatch.await(10, TimeUnit.SECONDS)).isTrue();
  7. assertThat(this.config.errorLatch.await(10, TimeUnit.SECONDS)).isTrue();
  8. assertThat(this.config.closeLatch.await(10, TimeUnit.SECONDS)).isTrue();
  9. MessageListenerContainer container = this.registry.getListenerContainer(CONTAINER_ID);
  10. assertThat(container.isRunning()).isFalse();
  11. InOrder inOrder = inOrder(this.consumer);
  12. inOrder.verify(this.consumer).subscribe(any(Collection.class), any(ConsumerRebalanceListener.class));
  13. inOrder.verify(this.consumer).poll(Duration.ofMillis(ContainerProperties.DEFAULT_POLL_TIMEOUT));
  14. inOrder.verify(this.consumer).commitSync(
  15. Collections.singletonMap(new TopicPartition("foo", 0), new OffsetAndMetadata(1L)));
  16. inOrder.verify(this.consumer).commitSync(
  17. Collections.singletonMap(new TopicPartition("foo", 0), new OffsetAndMetadata(2L)));
  18. inOrder.verify(this.consumer).commitSync(
  19. Collections.singletonMap(new TopicPartition("foo", 1), new OffsetAndMetadata(1L)));
  20. inOrder.verify(this.consumer).wakeup();
  21. inOrder.verify(this.consumer).unsubscribe();
  22. inOrder.verify(this.consumer).close();
  23. inOrder.verifyNoMoreInteractions();
  24. assertThat(this.config.count).isEqualTo(4);
  25. assertThat(this.config.contents.toArray()).isEqualTo(new String[]
  26. { "foo", "bar", "baz", "qux" });
  27. }

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

  1. /**
  2. * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
  3. */
  4. public void wakeup() {
  5. kafkaConsumer.wakeup();
  6. }

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

  1. @Override
  2. public void wakeup() {
  3. kafkaConsumer.wakeup();
  4. }
  5. }

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

  1. /**
  2. * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
  3. */
  4. public void wakeup() {
  5. kafkaConsumer.wakeup();
  6. }

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

  1. @Override
  2. public void wakeup() {
  3. consumer.wakeup();
  4. }
  5. }

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

  1. /**
  2. * Trigger the consumer's {@link KafkaConsumer#wakeup() wakeup()} method.
  3. */
  4. public void wakeup() {
  5. kafkaConsumer.wakeup();
  6. }

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

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

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

  1. public void shutdown() {
  2. closed.set( true );
  3. consumer.wakeup();
  4. }

相关文章