Index

A C D E F K M P Q R S 
All Classes and Interfaces|All Packages

A

aggregator - Variable in class me.kpavlov.finchly.queue.QueueSubscriber
 
awaitMessage(boolean, Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Same as MessageAggregator.awaitMessage(Duration, boolean, Predicate) using a default 5-second timeout.
awaitMessage(Duration, boolean, Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Blocks (polling, never sleeping the calling thread directly) until a message matching predicate arrives, or throws ConditionTimeoutException after timeout.
awaitMessage(Duration, Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Same as MessageAggregator.awaitMessage(Duration, boolean, Predicate) with extract defaulted to false (the message is not removed).
awaitMessage(Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Same as MessageAggregator.awaitMessage(Duration, boolean, Predicate) using a default 5-second timeout and extract defaulted to false.
awaitMessages(int, boolean, Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
awaitMessages(int, Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Same as MessageAggregator.awaitMessages(Duration, int, boolean, Predicate) using a default 5-second timeout and extract defaulted to false.
awaitMessages(Duration, int, boolean, Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Blocks until at least count messages matching predicate have arrived, or throws ConditionTimeoutException after timeout.
awaitMessages(Duration, int, Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Same as MessageAggregator.awaitMessages(Duration, int, boolean, Predicate) with extract defaulted to false (messages are not removed).

C

clear() - Method in class me.kpavlov.finchly.queue.MessageAggregator
Removes all stored messages.

D

deliver(T) - Method in class me.kpavlov.finchly.queue.QueueSubscriber
Called by subclasses for each message received from the transport.

E

extract(Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Atomically finds and removes the first message matching predicate.
extractAll(Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Atomically finds and removes all messages matching predicate, in arrival order.

F

find(Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Returns the first message matching predicate, without removing it.
findAll(Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Returns all messages matching predicate, in arrival order, without removing them.

K

KafkaQueuePublisher<K,T> - Class in me.kpavlov.finchly.kafka
QueuePublisher backed by a Kafka KafkaProducer, publishing every message to a single fixed topic and partition.
KafkaQueuePublisher(KafkaProducer<K, T>, String) - Constructor for class me.kpavlov.finchly.kafka.KafkaQueuePublisher
 
KafkaQueueSubscriber<K,T> - Class in me.kpavlov.finchly.kafka
QueueSubscriber backed by a Kafka Consumer (typically a KafkaConsumer).
KafkaQueueSubscriber(Consumer<K, T>, Collection<String>, MessageAggregator<T>) - Constructor for class me.kpavlov.finchly.kafka.KafkaQueueSubscriber
 

M

me.kpavlov.finchly.kafka - package me.kpavlov.finchly.kafka
 
me.kpavlov.finchly.queue - package me.kpavlov.finchly.queue
 
me.kpavlov.finchly.rabbitmq - package me.kpavlov.finchly.rabbitmq
 
MessageAggregator<T> - Class in me.kpavlov.finchly.queue
Thread-safe, order-retaining in-memory store for messages received by a QueueSubscriber.
MessageAggregator() - Constructor for class me.kpavlov.finchly.queue.MessageAggregator
 

P

publish(T) - Method in class me.kpavlov.finchly.kafka.KafkaQueuePublisher
 
publish(T) - Method in class me.kpavlov.finchly.queue.QueuePublisher
Publishes a single message.
publish(T) - Method in class me.kpavlov.finchly.rabbitmq.RabbitMqQueuePublisher
 
publishAll(Collection<T>) - Method in class me.kpavlov.finchly.queue.QueuePublisher
Publishes each message in order.
push(T) - Method in class me.kpavlov.finchly.queue.MessageAggregator
Appends a message, preserving arrival order.

Q

QueuePublisher<T> - Class in me.kpavlov.finchly.queue
Abstract base for a broker-specific publisher.
QueuePublisher() - Constructor for class me.kpavlov.finchly.queue.QueuePublisher
 
QueueSubscriber<T> - Class in me.kpavlov.finchly.queue
Abstract base for a broker-specific subscriber/receiver.
QueueSubscriber(MessageAggregator<T>) - Constructor for class me.kpavlov.finchly.queue.QueueSubscriber
 

R

RabbitMqQueuePublisher<T> - Class in me.kpavlov.finchly.rabbitmq
QueuePublisher backed by a RabbitMQ Channel, publishing every message (serialized via serializer) to a fixed exchange/routing key.
RabbitMqQueuePublisher(Channel, String, String, Function<T, byte[]>) - Constructor for class me.kpavlov.finchly.rabbitmq.RabbitMqQueuePublisher
 
RabbitMqQueueSubscriber<T> - Class in me.kpavlov.finchly.rabbitmq
QueueSubscriber backed by a RabbitMQ Channel, consuming from a single queue and delivering each message body (deserialized via deserializer) to the aggregator.
RabbitMqQueueSubscriber(Channel, String, Function<byte[], T>, MessageAggregator<T>) - Constructor for class me.kpavlov.finchly.rabbitmq.RabbitMqQueueSubscriber
 

S

size() - Method in class me.kpavlov.finchly.queue.MessageAggregator
Number of messages currently stored.
start() - Method in class me.kpavlov.finchly.kafka.KafkaQueueSubscriber
Subscribes and starts the poll loop.
start() - Method in class me.kpavlov.finchly.queue.QueueSubscriber
Starts listening for messages.
start() - Method in class me.kpavlov.finchly.rabbitmq.RabbitMqQueueSubscriber
Starts consuming, acknowledging each message only after it has been deserialized and delivered to the aggregator; a message that fails deserialization is negatively acknowledged (not requeued, to avoid a poison-message redelivery loop) rather than silently acknowledged and lost.
stop() - Method in class me.kpavlov.finchly.kafka.KafkaQueueSubscriber
Stops the poll loop and closes the consumer.
stop() - Method in class me.kpavlov.finchly.queue.QueueSubscriber
Stops listening and releases any transport resources (connections, threads).
stop() - Method in class me.kpavlov.finchly.rabbitmq.RabbitMqQueueSubscriber
Cancels the consumer.
A C D E F K M P Q R S 
All Classes and Interfaces|All Packages