Index
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
predicatearrives, or throwsConditionTimeoutExceptionaftertimeout. - awaitMessage(Duration, Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
-
Same as
MessageAggregator.awaitMessage(Duration, boolean, Predicate)withextractdefaulted tofalse(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 andextractdefaulted tofalse. - awaitMessages(int, boolean, Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
-
Same as
MessageAggregator.awaitMessages(Duration, int, boolean, Predicate)using a default 5-second timeout. - 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 andextractdefaulted tofalse. - awaitMessages(Duration, int, boolean, Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
-
Blocks until at least
countmessages matchingpredicatehave arrived, or throwsConditionTimeoutExceptionaftertimeout. - awaitMessages(Duration, int, Predicate<T>) - Method in class me.kpavlov.finchly.queue.MessageAggregator
-
Same as
MessageAggregator.awaitMessages(Duration, int, boolean, Predicate)withextractdefaulted tofalse(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 -
QueuePublisherbacked by a KafkaKafkaProducer, 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 - 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
-
QueuePublisherbacked by a RabbitMQChannel, publishing every message (serialized viaserializer) 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
-
QueueSubscriberbacked by a RabbitMQChannel, consuming from a single queue and delivering each message body (deserialized viadeserializer) 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.
All Classes and Interfaces|All Packages