org.springframework.cloud.stream.binder.ExtendedConsumerProperties.getInstanceIndex()方法的使用及代码示例

x33g5p2x  于2022-01-19 转载在 其他  
字(4.3k)|赞(0)|评价(0)|浏览(94)

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

ExtendedConsumerProperties.getInstanceIndex介绍

暂无

代码示例

代码示例来源:origin: spring-cloud/spring-cloud-stream-binder-kafka

public Collection<PartitionInfo> processTopic(final String group,
    final ExtendedConsumerProperties<KafkaConsumerProperties> extendedConsumerProperties,
    final ConsumerFactory<?, ?> consumerFactory, int partitionCount, boolean usingPatterns,
    boolean groupManagement, String topic) {
  Collection<PartitionInfo> listenedPartitions;
  Collection<PartitionInfo> allPartitions = usingPatterns ? Collections.emptyList()
      : getPartitionInfo(topic, extendedConsumerProperties, consumerFactory, partitionCount);
  if (groupManagement ||
      extendedConsumerProperties.getInstanceCount() == 1) {
    listenedPartitions = allPartitions;
  }
  else {
    listenedPartitions = new ArrayList<>();
    for (PartitionInfo partition : allPartitions) {
      // divide partitions across modules
      if ((partition.partition()
          % extendedConsumerProperties.getInstanceCount()) == extendedConsumerProperties
              .getInstanceIndex()) {
        listenedPartitions.add(partition);
      }
    }
  }
  this.topicsInUse.put(topic, new TopicInformation(group, listenedPartitions, usingPatterns));
  return listenedPartitions;
}

代码示例来源:origin: org.springframework.cloud/spring-cloud-stream-binder-kafka

public Collection<PartitionInfo> processTopic(final String group,
    final ExtendedConsumerProperties<KafkaConsumerProperties> extendedConsumerProperties,
    final ConsumerFactory<?, ?> consumerFactory, int partitionCount, boolean usingPatterns,
    boolean groupManagement, String topic) {
  Collection<PartitionInfo> listenedPartitions;
  Collection<PartitionInfo> allPartitions = usingPatterns ? Collections.emptyList()
      : getPartitionInfo(topic, extendedConsumerProperties, consumerFactory, partitionCount);
  if (groupManagement ||
      extendedConsumerProperties.getInstanceCount() == 1) {
    listenedPartitions = allPartitions;
  }
  else {
    listenedPartitions = new ArrayList<>();
    for (PartitionInfo partition : allPartitions) {
      // divide partitions across modules
      if ((partition.partition()
          % extendedConsumerProperties.getInstanceCount()) == extendedConsumerProperties
              .getInstanceIndex()) {
        listenedPartitions.add(partition);
      }
    }
  }
  this.topicsInUse.put(topic, new TopicInformation(group, listenedPartitions, usingPatterns));
  return listenedPartitions;
}

代码示例来源:origin: org.springframework.cloud/spring-cloud-stream-binder-rabbit-core

private Binding declareConsumerBindings(String name, ExtendedConsumerProperties<RabbitConsumerProperties> properties,
                  Exchange exchange, boolean partitioned, Queue queue) {
  if (partitioned) {
    return partitionedBinding(name, exchange, queue, properties.getExtension(), properties.getInstanceIndex());
  }
  else {
    return notPartitionedBinding(exchange, queue, properties.getExtension());
  }
}

代码示例来源:origin: spring-cloud/spring-cloud-stream-binder-rabbit

private Binding declareConsumerBindings(String name, ExtendedConsumerProperties<RabbitConsumerProperties> properties,
                  Exchange exchange, boolean partitioned, Queue queue) {
  if (partitioned) {
    return partitionedBinding(name, exchange, queue, properties.getExtension(), properties.getInstanceIndex());
  }
  else {
    return notPartitionedBinding(exchange, queue, properties.getExtension());
  }
}

代码示例来源:origin: org.springframework.cloud/spring-cloud-stream-binder-kinesis

for (int i = 0; i < shards.size(); i++) {
  if ((i % properties.getInstanceCount()) == properties.getInstanceIndex()) {
    KinesisShardOffset shardOffset = new KinesisShardOffset(kinesisShardOffset);
    shardOffset.setStream(destination.getName());

代码示例来源:origin: spring-cloud/spring-cloud-stream-binder-aws-kinesis

.getInstanceIndex()) {
KinesisShardOffset shardOffset = new KinesisShardOffset(
    kinesisShardOffset);

代码示例来源:origin: org.springframework.cloud/spring-cloud-stream-binder-kafka11

.getInstanceIndex()) {
listenedPartitions.add(partition);

代码示例来源:origin: org.springframework.cloud/spring-cloud-stream-binder-rabbit-core

String partitionSuffix = "-" + properties.getInstanceIndex();
queueName += partitionSuffix;

代码示例来源:origin: spring-cloud/spring-cloud-stream-binder-rabbit

String partitionSuffix = "-" + properties.getInstanceIndex();
queueName += partitionSuffix;

相关文章