Pubblica messaggi su un argomento in un modello Pub/Sub
Il modello Pub/Sub (Editore/Abbonato) in NCache consente la messaggistica disaccoppiata e basata sugli eventi. Gli editori possono inviare messaggi a un argomento utilizzando ITopic interfaccia senza sapere chi sono gli abbonati. Fornisce anche registrazioni di eventi per errori di recapito dei messaggi, ricezione di messaggi ed eliminazione di argomenti. Fornisce Publish metodo, che pubblica il messaggio su un argomento specifico nella cache. Durante la pubblicazione dei messaggi su un argomento, l'editore può impostare il opzione di consegna per i messaggi, scadenza del messaggio, notifica di mancata consegna del messaggioe Notifica di eliminazione dell'argomentoQui descriviamo come pubblicare messaggi, pubblicare i messaggi in modo asincrono, pubblicare messaggi in bloccoe pubblicare messaggi ordinati.
Prerequisiti
Prima di usare il NCache Per le API lato client, assicurarsi che siano soddisfatti i seguenti prerequisiti:
- Installa i seguenti pacchetti NuGet nella tua applicazione client .NET:
- Includere i seguenti spazi dei nomi nell'applicazione:
- La cache deve essere in esecuzione.
- Assicurati che i dati aggiunti lo siano serializzabile.
- Per i dettagli dell'API, fare riferimento a: ICache, CacheItem, Argomento, Pubblica, GetTopic, PubblicaAsync, Data di scadenza, Errore di consegna del messaggio, OnTopic Eliminato, Opzione di consegna, Messaggio, Pubblica in blocco, MessageFailedEventArgs, ArgomentoDeleteEventArgs, Servizio di messaggistica.
- Aggiungi le seguenti dipendenze Maven per la tua applicazione client Java in
pom.xml file:
<dependency>
<groupId>com.alachisoft.ncache</groupId>
<!--for NCache Enterprise-->
<artifactId>ncache-client</artifactId>
<version>x.x.x</version>
</dependency>
- Importa i seguenti pacchetti nell'applicazione client Java:
- La cache deve essere in esecuzione.
- Assicurati che i dati aggiunti lo siano serializzabile.
- Per i dettagli dell'API, fare riferimento a: Cache, CacheItem, Argomento, getTopic, pubblicare, Messaggio, setExpirationTime, Ascoltatore di argomenti, addMessageDeliveryFailureListener, addTopicDeletedListener, Opzione di consegna, Tutti, pubblicareAsync, pubblicare Bulk, ArgomentoDeleteEventArgs, MessageFailedEventArgs, MessageEventArgs, getMessagingService.
- Installa i seguenti pacchetti nella tua applicazione client Python:
- Importa i seguenti pacchetti nella tua applicazione:
- La cache deve essere in esecuzione.
- Assicurati che i dati aggiunti lo siano serializzabile.
- Per i dettagli dell'API, fare riferimento a: Cache, CacheItem, get_topic, get_messaging_service, Intervallo di tempo, imposta_scadenza_time, add_message_delivery_failure_listener, add_topic_deleted_listener, pubblica_async, pubblica_in blocco, MessageEventArgs, get_nome_argomento, MessageFailedEventArgs, get_message_failure_reason, ArgomentoDeleteEventArgs, Argomento, Messaggio, Opzione di consegna, pubblicare, get_nome_argomento.
- Installa uno dei seguenti pacchetti NuGet nella tua applicazione client .NET:
- Enterprise:
Install-Package Alachisoft.NCache.SDK -Version 4.9.1.0
- Crea una nuova applicazione console.
- Assicurati che i dati aggiunti lo siano serializzabile.
- Aggiungi NCache Referenze individuando
%NCHOME%\NCache\bin\assembly\4.0 e aggiungendo Alachisoft.NCache.Web and Alachisoft.NCache.Runtime come appropriato.
- Includi il
Alachisoft.NCache.Web.Caching and Alachisoft.NCache.Runtime.Caching spazi dei nomi nella tua applicazione.
- Per saperne di più sul NCache API legacy, scaricala NCache 4.9 documenti disponibili come file .zip file sul Alachisoft Sito web.
Pubblica messaggi
Il seguente esempio di codice mostra come:
- Crea argomenti dedicati per i messaggi relativi agli ordini.
- Registrati
MessageDeliveryFailure evento per l'argomento.
- Registrati
OnTopicDeleted evento per l'argomento.
- Crea messaggi per ogni argomento, abilitando le opzioni di scadenza e consegna.
- Pubblica i messaggi.
try
{
// Precondition: Cache is already connected
// Topic "orderTopic" exists in cache
string topicName = "orderTopic";
// Get the Topic
ITopic orderTopic = cache.MessagingService.GetTopic(topicName);
if (orderTopic != null)
{
// Create the object to be sent in message
Order order = FetchOrderFromDB(10248);
// Create the new message with the object order
var orderMessage = new Message(order);
// Set the expiration time of the message
orderMessage.ExpirationTime = TimeSpan.FromSeconds(5000);
// Register message delivery failure
orderTopic.MessageDeliveryFailure += OnFailureMessageReceived;
// Register Topic deletion notification
orderTopic.OnTopicDeleted = TopicDeleted;
// Publish the order with delivery option set as all and register message delivery failure
orderTopic.Publish(orderMessage, DeliveryOption.All, true);
}
else
{
// No Topic exists
}
}
catch (OperationFailedException ex)
{
if (ex.ErrorCode == NCacheErrorCodes.MESSAGE_ID_ALREADY_EXISTS)
{
// Message ID already exists, specify a new ID
}
if (ex.ErrorCode == NCacheErrorCodes.TOPIC_DISPOSED)
{
// Specified topic has been disposed
}
if (ex.ErrorCode == NCacheErrorCodes.PATTERN_BASED_PUBLISHING_NOT_ALLOWED)
{
// Message publishing on pattern based topic is not allowed
// Get non-pattern based topic
}
else
{
// Exception can occur due to:
// Connection Failures
// Operation Timeout
// Operation performed during state transfer
}
}
catch (Exception ex)
{
// Any other generic exception like ArgumentNullException or ArgumentException
// Topic name is null/empty
}
try
{
// Precondition: Cache is already connected
// Already existing Topic
String topicName = "orderTopic";
// Get Topic
Topic orderTopic = cache.getMessagingService().getTopic(topicName);
if (topicName != null) {
// Create object to be sent in the message
Order order = fetchOrdersFromDB(1100);
// Create new message
Message orderMessage = new Message(order);
TimeSpan expiryTime = new TimeSpan(12, 12, 12);
// Set expiration time of the message
orderMessage.setExpirationTime(expiryTime);
// Register message delivery failure
MyTopicListener topicListener = new MyTopicListener();
orderTopic.addMessageDeliveryFailureListener(topicListener);
orderTopic.addTopicDeletedListener(topicListener);
// Publish the order with delivery option set as All and register message delivery failure
orderTopic.publish(orderMessage, DeliveryOption.All, true);
} else
{
// No Topic exists
}
}
catch (OperationFailedException exception)
{
if (exception.getErrorCode() == NCacheErrorCodes.MESSAGE_ID_ALREADY_EXISTS) {
// Message ID already exists. Specify a new ID
}
if (exception.getErrorCode() == NCacheErrorCodes.TOPIC_DISPOSED) {
// Specified topic has been disposed
}
if (exception.getErrorCode() == NCacheErrorCodes.PATTERN_BASED_PUBLISHING_NOT_ALLOWED) {
// Message publishing on pattern based topic is not allowed
// Get non-pattern based topic
} else {
// Exception can occur due to:
// Connection Failures
// Operation Timeout
// Operation performed during state transfer
}
}
catch (Exception exception)
{
// Any generic exception like IllegalArgumentException or NullPointerException
}
try:
# Precondition: Cache is already connected
# Topic "orderTopic" exists in the cache
topic_name = "orderTopic"
# Get the Topic
order_topic = cache.get_messaging_service().get_topic(topic_name)
if order_topic is not None:
# Create object to be sent in message
order = fetch_order_from_db(10248)
# Create a new message with the object order
order_message = Message(order)
# Set the expiration time of the message
expiry_time = TimeSpan(0, 12, 12, 12, 0)
order_message.set_expiration_time(expiry_time)
# Register message delivery failure
order_topic.add_message_delivery_failure_listener(on_message_delivery_failure)
# Register Topic deletion notification
order_topic.add_topic_deleted_listener(on_topic_deleted)
# Publish the order with delivery option set as All and register message delivery failure
order_topic.publish(order_message, DeliveryOption.ALL, True)
else:
# No Topic exists
except Exception as ex:
# Exception can occur due to:
# - Connection failures
# - Operation timeout
# - Topic does not exist / invalid arguments
print("Operation failed: " + str(ex))
try
{
// This is an async method
// Precondition: Cache is already connected
// Topic "orderTopic" exists in cache
let topicName = "orderTopic";
// Get the Topic
let messagingService = await cache.getMessagingService();
let orderTopic = await messagingService.getTopic(topicName);
if (orderTopic !== null)
{
// Create the object to be sent in message
let order = await FetchOrderFromDB(10248);
// Create the new message with the object order
let orderMessage = new ncache.Message(order);
// Set the expiration time of the message
let expiryTime = new ncache.TimeSpan(0, 0, 0, 5000);
orderMessage.setExpirationTime(expiryTime);
// Register message delivery failure
orderTopic.addMessageDeliveryFailureListener(topicListener);
// Register Topic deletion notification
orderTopic.addTopicDeletedListener(topicListener);
// Publish the order with delivery option set as all and register message delivery failure
await orderTopic.publish(orderMessage, ncache.DeliveryOption.All, null, true);
}
else
{
// No Topic exists
}
}
catch (error) {
// Handle any errors
}
try
{
// Using NCache Enterprise 4.9.1
// Precondition: Cache is already connected
// Topic "orderTopic" exists in cache
string topicName = "orderTopic";
// Get the Topic
ITopic orderTopic = cache.MessagingService.CreateTopic(topicName);
if (orderTopic != null)
{
// Create the object to be sent in message
Order order = FetchOrderFromDB(10248);
// Create the new message with the object order
Message orderMessage = new Message(order);
// Set the expiration time of the message
orderMessage.ExpirationTime = new TimeSpan(0, 0, 150);
// Register message delivery failure
orderTopic.MessageDeliveryFailure += FailureMessageReceived;
// Register Topic deletion notification
orderTopic.OnTopicDeleted = TopicDeleted;
// Publish the order with delivery option set as all and register message delivery failure
orderTopic.Publish(orderMessage, DeliveryOption.All, true);
}
else
{
// No Topic exists
}
}
catch (OperationFailedException ex)
{
if (ex.ErrorCode == NCacheErrorCodes.MESSAGE_ID_ALREADY_EXISTS)
{
// Message ID already exists, specify a new ID
}
if (ex.ErrorCode == NCacheErrorCodes.TOPIC_DISPOSED)
{
// Specified topic has been disposed
}
if (ex.ErrorCode == NCacheErrorCodes.PATTERN_BASED_PUBLISHING_NOT_ALLOWED)
{
// Message publishing on pattern based topic is not allowed
// Get non-pattern based topic
}
else
{
// Exception can occur due to:
// Connection Failures
// Operation Timeout
// Operation performed during state transfer
}
}
catch (Exception ex)
{
// Any other generic exception like ArgumentNullException or ArgumentException
// Topic name is null/empty
}
Note:
Per garantire che l'operazione sia a prova di errore, si consiglia di gestire eventuali potenziali eccezioni all'interno dell'applicazione, come spiegato in Gestione dei guasti.
Pubblica in modo asincrono
I messaggi possono essere pubblicati sull'argomento in modo asincrono utilizzando PubblicaAsync In questo modo, l'applicazione non attende il completamento dell'operazione per eseguire la successiva. Il controllo viene restituito immediatamente, consentendo all'applicazione di continuare l'elaborazione. L'esempio seguente mostra come pubblicare un messaggio in modo asincrono.
// Precondition: Cache is already connected
// Topic "orderTopic" exists in cache
string topicName = "orderTopic";
// Get the Topic
ITopic orderTopic = cache.MessagingService.GetTopic(topicName);
if (orderTopic != null)
{
// Create the object to be sent in message
Order order = FetchOrderFromDB(10248);
// Create the new message with the object order
var orderMessage = new Message(order);
// Set the expiration time of the message
orderMessage.ExpirationTime = TimeSpan.FromSeconds(5000);
// Register message delivery failure
orderTopic.MessageDeliveryFailure += OnFailureMessageReceived;
// Register Topic deletion notification
orderTopic.OnTopicDeleted = TopicDeleted;
// Publish the order with delivery option set as all and register message delivery failure
Task task = orderTopic.PublishAsync(orderMessage, DeliveryOption.All, true);
if(task.IsFaulted)
{
// Task Failed
}
}
else
{
// No Topic exists
}
// Precondition: Cache is already connected
// Topic orderTopic already exists in the cache
String topicName = "orderTopic";
// Get the Topic
Topic orderTopic = cache.getMessagingService().getTopic(topicName);
if (orderTopic != null) {
// Create object to be sent in the message
Order order = fetchOrdersFromDB(1100);
// Create new message
Message orderMessage = new Message(order);
TimeSpan expiryTime = new TimeSpan(12, 12, 12);
// Set expiration time of the message
orderMessage.setExpirationTime(expiryTime);
// Register message delivery failure
MyTopicListener topicListener = new MyTopicListener();
orderTopic.addMessageDeliveryFailureListener(topicListener);
orderTopic.addTopicDeletedListener(topicListener);
// Publish the order with delivery option set as All and register message delivery failure
TimeScheduler.Task task = (TimeScheduler.Task) orderTopic.publishAsync(orderMessage, DeliveryOption.All, true);
if (task.IsCancelled()) {
// Task cancelled
}
} else {
// No Topic exists
}
# Precondition: Cache is already connected
# Topic "orderTopic" exists in the cache
topic_name = "orderTopic"
# Get the Topic
order_topic = cache.get_messaging_service().get_topic(topic_name)
if order_topic is not None:
# Create object to be sent in message
order = fetch_order_from_db(10248)
# Create a new message with the object order
order_message = Message(order)
# Set the expiration time of the message
expiry_time = TimeSpan(0, 1, 23, 20, 0)
order_message.set_expiration_time(expiry_time)
# Register message delivery failure
order_topic.add_message_delivery_failure_listener(on_message_delivery_failure)
# Register Topic deletion notification
order_topic.add_topic_deleted_listener(on_topic_deleted)
# Publish the order with delivery option set as All and register message delivery failure
try:
await order_topic.publish_async(order_message, DeliveryOption.ALL, True)
print("Message published asynchronously.")
except Exception as e:
print(f"Async publish failed: {e}")
else:
# No Topic exists
Pubblica messaggi collettivi
È possibile pubblicare più messaggi in una singola chiamata utilizzando il Pubblica in blocco metodo. Ciò migliora le prestazioni e l'utilizzo della memoria poiché un sacco di messaggi verrà combinato e pubblicato in una singola chiamata. Il codice seguente prende un'istanza di un argomento già creato ordineArgomentoe mostra la pubblicazione in blocco dei messaggi nell'argomento.
// Precondition: Cache is already connected
// Topic "orderTopic" exists in cache
ITopic topic = cache.MessagingService.GetTopic("orderTopic");
if (topic != null)
{
// Create dictionary for storing bulk
List<Tuple<Message, DeliveryOption>> messageList = new List<Tuple<Message, DeliveryOption>>();
Order[] orders = FetchOrdersFromDB();
for (int i = 0; i < 100; i++)
{
Message message = new Message(orders[i]);
message.ExpirationTime = TimeSpan.FromSeconds(10000);
messageList.Add(new Tuple<Message, DeliveryOption>(message, DeliveryOption.All));
}
// Register message delivery failure
topic.MessageDeliveryFailure += OnFailureMessageReceived;
// Register Topic deletion notification
topic.OnTopicDeleted = TopicDeleted;
// Publish the order with delivery option set as all and register message delivery failure
// In case of failed publishing of messages, exceptions will be returned
IDictionary<Message, Exception> keys = topic.PublishBulk(messageList, true);
}
// Precondition: Cache is already connected
// Topic already exists
String topicName = "orderTopic";
String customerID = "DUMON";
Topic topic = cache.getMessagingService().getTopic(topicName);
if (topic != null) {
// Create dictionary for storing bulk
Map messageMap = new HashMap();
Order[] orders = fetchOrdersFromDb(customerID);
for (int i = 0; i < 100; i++) {
Message message = new Message(orders[i]);
message.setExpirationTime(TimeSpan.FromSeconds(10000));
messageMap.put(message, DeliveryOption.All);
}
MyTopicListener topicListener = new MyTopicListener();
// Register message delivery failure
topic.addMessageDeliveryFailureListener(topicListener);
// Register Topic deletion notification
topic.addTopicDeletedListener(topicListener);
Map<Message, Exception> keys = topic.publishBulk(messageMap, true);
} else
{
// Topic is null
}
# Precondition: Cache is already connected
# Topic "orderTopic" exists in the cache
topic_name = "orderTopic"
# Get the Topic
topic = cache.get_messaging_service().get_topic(topic_name)
if topic is not None:
# Create Dict for storing messages in bulk
messages_map = {}
orders = fetch_orders_from_db()
for i in range(0, 100):
message = Message(orders[++i])
message.set_expiration_time(TimeSpan.from_seconds(10000))
messages_map[message] = DeliveryOption.ALL
# Register message delivery failure
topic.add_message_delivery_failure_listener(on_message_delivery_failure)
# Register Topic deletion notification
topic.add_topic_deleted_listener(on_topic_deleted)
# Publish the order with delivery option set as All and register message delivery failure
# In case of failed publishing of messages, exception will be returned
keys = topic.publish_bulk(messages_map, True)
// Precondition: Cache is already connected
// This is an async method
// Topic "orderTopic" exists in the cache
let messagingService = await cache.getMessagingService();
let topic = await messagingService.getTopic("orderTopic");
if (topic != null)
{
// Create Map for storing messages in bulk
let messageMap = new Map();
const orders = FetchOrdersFromDB();
for (var i = 0; i < 100; i++)
{
let message = new ncache.Message(orders[i]);
message.setExpirationTime(new ncache.TimeSpan(0, 0, 0, 10000));
messageMap.set(message, ncache.DeliveryOption.All);
}
// Register message delivery failure
topic.addMessageDeliveryFailureListener(topicListener);
// Register Topic deletion notification
topic.addTopicDeletedListener(topicListener);
// Publish the order with delivery option set as All and register message delivery failure
// In case of failed publishing of messages, exception will be returned
let keys = topic.publishBulk(messageMap, true);
}
else
{
// No Topic exists
}
Pubblica messaggi ordinati
Note:
Questa funzione è disponibile solo in NCache Dal 5.2 in poi.
I messaggi possono essere pubblicati specificando un nome di sequenza che fa sì che i messaggi vengano pubblicati in un ordine specifico. Per specificare messaggi ordinati, un nome sequenza stringa viene aggiunto alla catena dei messaggi che assicura la pubblicazione di tutti i messaggi appartenenti a un nome sequenza specifico sullo stesso nodo server. Nell'esempio riportato di seguito, un nome sequenza viene aggiunto con i messaggi e i messaggi vengono quindi pubblicati utilizzando Publish metodo.
// Precondition: Cache is already connected
// Specify the Topic name that already exists
string topicName = "orderTopic";
// Get the Topic with the specified name
ITopic orderTopic = cache.MessagingService.GetTopic(topicName);
if (orderTopic != null)
{
for (int i = 0; i < 30; i++)
{
// Create the object to be sent in message
Order order = FetchOrderFromDB(10248);
// Create the new message with the object order
var orderMessage = new Message(order);
// Specify a unique sequence name for the messages
string sequenceName = "OrderMessages";
// Set the expiration time of the message
orderMessage.ExpirationTime = TimeSpan.FromSeconds(5000);
// Publish message with the sequence name
orderTopic.Publish(orderMessage, DeliveryOption.All, sequenceName, true);
}
}
else
{
// No Topic found
}
// Precondition: Cache is already connected
// Specify the Topic name that already exists
String topicName = "orderTopic";
// Get the Topic with the specified name
Topic orderTopic = cache.getMessagingService().getTopic(topicName);
if (orderTopic != null)
{
for (int i = 0; i < 30; i++) {
// Create the object to be sent in the message
Order order = fetchOrdersFromDB(10248);
// Create the new message with the object order
var orderMessage = new Message(order);
// Specify a unique sequence name for the message
String sequenceName = "OrderMessage";
// Set the expiration time of the message
orderMessage.setExpirationTime(TimeSpan.FromSeconds(5000));
// Publish message with the sequence name
orderTopic.publish(orderMessage, DeliveryOption.All, sequenceName, true);
}
} else {
// No Topic found
}
# Precondition: Cache is already connected
# Specify the Topic name that already exists
topic_name = "orderTopic"
# Get the Topic with the specified name
order_topic = cache.get_messaging_service().get_topic(topic_name)
if topic_name is not None:
for i in range(0, 30):
# Create the object to be sent in the message
order = fetch_order_from_db(10248)
# Create the new message with the object order
order_message = Message(order)
# Specify a unique sequence name for the message
sequence_name = "OrderMessages"
# Set the expiration time of the message
order_message.set_expiration_time(TimeSpan.from_seconds(5000))
# Publish message with the sequence name
order_topic.publish(order_message, DeliveryOption.ALL, True, sequence_name)
else:
# No Topic found
// Precondition: Cache is already connected
// This is an async method
// Specify the Topic name that already exists
let topicName = "orderTopic";
// Get the Topic with the specified name
let messagingService = await cache.getMessagingService();
let orderTopic = await messagingService.getTopic(topicName);
if (orderTopic != null)
{
for (var i = 0; i < 30; i++)
{
// Create the object to be sent in message
let order = await FetchOrderFromDB();
// Create the new message with the object order
let orderMessage = new ncache.Message(order);
// Specify a unique sequence name for the message
let sequenceName = "OrderMessage";
// Set the expiration time of the message
orderMessage.setExpirationTime(new ncache.TimeSpan(0, 0, 0, 5000));
// Publish message with the sequence name
orderTopic.publish(
orderMessage,
ncache.DeliveryOption.All,
sequenceName,
true
);
}
}
else
{
// No Topic found
}
Registra le richiamate
private void OnFailureMessageReceived(object sender, MessageFailedEventArgs args)
{
// Failure reason can be get from args.MessageFailureReason
}
private void TopicDeleted(object sender, TopicDeleteEventArgs args)
{
// Deleted Topic is args.TopicName
}
@Override
public void onTopicDeleted(Object sender, TopicDeleteEventArgs args) {
// Deleted Topic is args.getTopicName();
}
@Override
public void onMessageDeliveryFailure(Object sender, MessageFailedEventArgs args) {
// Failure reason is args.getMessageFailureReason();
}
@Override
public void onMessageReceived(Object sender, MessageEventArgs args) {
// Perform operations
}
def on_topic_deleted(sender: object, args: ncache.TopicDeleteEventArgs):
# Perform operations
print("Deleted Topic is " + args.get_topic_name())
def on_message_delivery_failure(sender: object, args: ncache.MessageFailedEventArgs):
# Perform operations
print("Failure reason is " + str(args.get_message_failure_reason()))
def on_message_received(sender: object, args: ncache.MessageEventArgs):
# Perform operations
print("Message received from Topic " + args.get_topic_name())
function onMessageDeliveryFailure(sender, args)
{
// Failure reason is args.getMessageFailureReason()
}
function onTopicDeleted(sender, args)
{
// Deleted Topic is args.getTopicName()
}
// Using NCache Enterprise 4.9.1
private void FailureMessageReceived(object sender, MessageFailedEventArgs args)
{
// Failure reason can be get from args.MessageFailureReason
}
private void TopicDeleted(object sender, TopicDeleteEventArgs args)
{
// Deleted Topic is args.TopicName
}
Risorse addizionali
NCache fornisce un'applicazione di esempio per Pub/Sub su GitHub.
Vedere anche
.NETTO: Alachisoft.NCache.Memorizzazione.della.cache spazio dei nomi.
Giava: com.alachisoft.ncache.runtime.caching pacchetto.
Pitone: servizi.ncache.client modulo.
Node.js: Argomento classe.