本文整理了Java中org.springframework.cloud.stream.binder.ExtendedConsumerProperties.getInstanceIndex()
方法的一些代码示例,展示了ExtendedConsumerProperties.getInstanceIndex()
的具体用法。这些代码示例主要来源于Github
/Stackoverflow
/Maven
等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。ExtendedConsumerProperties.getInstanceIndex()
方法的具体详情如下:
包路径:org.springframework.cloud.stream.binder.ExtendedConsumerProperties
类名称: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;
内容来源于网络,如有侵权,请联系作者删除!