Messaggistica Pub/Sub e flussi

Molte applicazioni server oggi necessitano di coordinare il lavoro con altre applicazioni e condividere messaggi e dati in modo asincrono. E spesso ciò avviene in un ambiente ad alta transazione. Per ambienti con transazioni elevate, NCache è un candidato ideale per questi comunicazioni asincrone guidate da eventi.

NCache fornisce un ricco set di funzionalità di messaggistica e flussi per soddisfare le esigenze di queste applicazioni, incluse Messaggistica Pub/Sub, Query continue, Eventi a livello di oggettoe Eventi a livello di cache. Ciascuno è descritto di seguito.

NCache Architettura di messaggistica Pub/Sub per .NET, Java e Python
NCache Architettura di messaggistica che dimostra scalabilità lineare per applicazioni distribuite.
 

Messaggistica Pub/Sub

NCache fornisce un potente Messaggistica Pub/Sub piattaforma per applicazioni in cui le applicazioni Editore possono inviare messaggi senza sapere chi sono gli Abbonati. I messaggi possono essere classificati in Argomenti quindi gli abbonati possono iscriversi ad argomenti specifici e ricevere questi messaggi in modo asincrono.

  • - Editore: Pubblica il messaggio affinché possa essere fruito da più abbonati. L'editore non sa chi sono gli abbonati e quanti sono.
  • - Subscriber: Si registra per ricevere messaggi relativi a un argomento specifico. L'abbonato non sa chi è l'editore né quanti altri abbonati ci sono.
  • - Argomento: Tutti i messaggi vengono pubblicati e ricevuti per ogni argomento. Gli editori pubblicano messaggi per quell'argomento e gli abbonati ricevono messaggi da quell'argomento. Puoi avere più argomenti in NCache.
  • - Messaggio: Questo è il messaggio effettivo inviato dai Publisher e utilizzato dai Subscriber. Supporta tutti gli oggetti serializzabili, inclusi dati binari, JSON, XML e tipi primitivi.

Si può facilmente uso NCache Messaggistica Pub/Sub per applicazioni distribuite una volta compresi i componenti della messaggistica Pub/Sub e il loro funzionamento.

NCache Diagramma Pub/Sub che mostra i Publisher che inviano messaggi agli Argomenti e i Subscriber che li ricevono.
NCache Modello Pub/Sub che mostra la comunicazione asincrona e disaccoppiata tramite argomenti.

Di seguito è riportato un esempio di codice sorgente su come un editore pubblica un messaggio ad un argomento e poi a l'abbonato si iscrive all'argomento.

ITopic topic = _cache.MessagingService.CreateTopic(topicName)

// Creating new message containing a “payload” 
var message = new Message(payload);

// Set the expiration TimeSpan of the message
message.ExpirationTime = TimeSpan.FromSeconds(5000);

// Publishing the above created message.
topic.Publish(message, DeliveryOption.All, true);

// Get the already created topic to subscribe to it
ITopic orderTopic = cache.MessagingService.GetTopic(topic);

// Create and register subscribers for the topic created
// MessageRecieved callback is registered
orderTopic.CreateSubscription(MessageReceived);
String topicName = "orderUpdates";
Topic topic = cache.getMessagingService().createTopic(topicName);

// Creating new message containing a “payload”
Customer payload = new Customer("David Jones", "XYZCompany");
var message = new Message(payload);

// Set the expiration TimeSpan of the message
message.setExpirationTime(TimeSpan.FromSeconds(5000));

// Publishing the above created message.
topic.publish(message, DeliveryOption.All, true);

// Get the already created topic to subscribe to it
Topic orderTopic = cache.getMessagingService().getTopic(topicName);

// Create and register subscribers for the topic created
orderTopic.createSubscription(new MessageReceivedListener() {
    @Override
     public void onMessageReceived(Object sender, MessageEventArgs args) {
// message received
     }
});
 

Interrogazione continua (CQ)

Interrogazione continua è una caratteristica potente di NCache che consente di monitorare le modifiche a un "set di dati" nella cache distribuita in base a query di tipo SQL. Crei una query continua (CQ) con criteri simili a SQL rispetto alla cache e la registri NCacheIn questo modo si elimina la necessità di effettuare il polling del database, riducendo il carico del database e garantendo che l'applicazione visualizzi sempre il set di dati più aggiornato in tempo reale. NCache poi inizia a monitorare eventuali modifiche a questo "set di dati" nella cache tra cui:

  • - Aggiungere: se un elemento viene aggiunto alla cache che corrisponde ai criteri CQ, il client viene avvisato insieme a questo oggetto.
  • - Aggiornare: se nella cache viene aggiornato un elemento che corrisponde ai criteri CQ, il client viene avvisato insieme a questo oggetto.
  • - Rimuovere: se un elemento viene rimosso dalla cache che corrisponde ai criteri CQ, il client viene avvisato insieme a questo oggetto.

Ecco come puoi creare una query continua con criteri specifici e registrarti sul server.

// Query for required operation
string query = "SELECT $VALUE$ FROM FQN.Customer WHERE Country = ?";

var queryCommand = new QueryCommand(query);
queryCommand.Parameters.Add("Country", "USA");

// Create Continuous Query
var cQuery = new ContinuousQuery(queryCommand);

// Register to be notified when a qualified item is added to the cache
cQuery.RegisterNotification(new QueryDataNotificationCallback(QueryItemCallBack), EventType.ItemAdded, EventDataFilter.None);

// Register continuousQuery on server 
cache.MessagingService.RegisterCQ(cQuery);
String cacheName = "demoCache";

// Initialize an instance of the cache to begin performing operations:
Cache cache = CacheManager.getCache(cacheName);

// Query for required operation
String query = "SELECT $VALUE$ FROM com.alachisoft.ncache.sample.Customer WHERE country = ?";

var queryCommand = new QueryCommand(query);
queryCommand.getParameters().put("Country", "USA");

// Create Continuous Query
var cQuery = new ContinuousQuery(queryCommand);

// Register to be notified when a qualified item is added to the cache
cQuery.addDataModificationListener(new QueryDataModificationListener() {
	@Override
	public void onQueryDataModified(String key, CQEventArg args) {
		switch (args.getEventType())
		{
			case ItemAdded:
				// 'key' has been added to the cache
				break;
		}
	}
  }, EnumSet.allOf(EventType.class), EventDataFilter.None);

// Register continuousQuery on server
cache.getMessagingService().registerCQ(cQuery);

Per ulteriori dettagli, puoi vedere come utilizzare la query continua nei documenti della cache.

 

Notifiche di eventi

I clienti si registrano con NCache ricevere notifiche di eventi sono interessati fornendo una funzione di callback. NCache fornisce filtri con eventi utilizzati per specificare la quantità di informazioni restituite con l'evento. Questi filtri sono Nessuno (predefinito), Metadati and Dati con metadati.

Sono supportati i seguenti tipi di eventi:

  • - Eventi a livello di oggetto: il client si è registrato per ricevere una notifica ogni volta che uno specifico elemento memorizzato nella cache (basato su una chiave) viene aggiornato o rimosso dalla cache e fornisce una richiamata per esso. Quando questo elemento viene aggiornato o rimosso, il client viene avvisato.
  • - Eventi a livello di cache: Gli eventi a livello di cache sono gli eventi generali che possono essere utili quando tutte le applicazioni client devono essere informate di qualsiasi modifica nella cache da parte di qualsiasi applicazione client connessa.

Il codice fornito di seguito mostra la registrazione degli eventi a livello di articolo.

public static void RegisterItemLevelEvents()
{
	// Get cache
	string cacheName = "demoCache";
	var cache = CacheManager.GetCache(cacheName); 

	// Register Item level events
	string key = "Product:1001";
	cache.MessagingService.RegisterCacheNotification(key, CacheDataNotificationCallbackImpl, EventType.ItemUpdated | EventType.ItemUpdated | EventType.ItemRemoved, EventDataFilter.DataWithMetadata);
}
public static void CacheDataNotificationCallbackImpl(string key, CacheEventArg cacheEventArgs)
{
	switch (cacheEventArgs.EventType)
	{
	   case EventType.ItemAdded:
		   // 'key' has been added to the cache
		   break;
	   case EventType.ItemUpdated:
		   // 'key' has been updated in the cache
		   break;
	   case EventType.ItemRemoved:
		   // 'key' has been removed from the cache
		   break;
	}
}
private static void RegisterItemLevelEvents() throws Exception {
        String cacheName = "demoCache";

        // Initialize an instance of the cache to begin performing operations:
        Cache cache = CacheManager.getCache(cacheName);
        String key = "Product:1001";
        cache.getMessagingService().addCacheNotificationListener(key, new CacheDataModificationListenerImpl(), EnumSet.allOf(EventType.class), EventDataFilter.DataWithMetadata);
    }
private static class CacheDataModificationListenerImpl implements CacheDataModificationListener {
	@Override
	public void onCacheDataModified(String key, CacheEventArg eventArgs) {
		switch (eventArgs.getEventType())
		{
			case ItemAdded:
				// 'key' has been added to the cache
				break;
			case ItemUpdated:
				// 'key' has been updated in the cache
				break;
			case ItemRemoved:
				// 'key' has been removed from the cache
				break;
		}
	}
	@Override
	public void onCacheCleared(String cacheName) {
		//cacheName cleared
	}
}

Per gli eventi a livello di cache, puoi utilizzare lo stesso callback utilizzato negli eventi a livello di elemento.

public static void RegisterCacheLevelEvents()
{
	// Get cache
	string cacheName = "demoCache";
	var cache = CacheManager.GetCache(cacheName);

	// Register cache level events
	cache.MessagingService.RegisterCacheNotification(EventsSample.CacheDataNotificationCallbackImpl, EventType.ItemUpdated | EventType.ItemUpdated | EventType.ItemRemoved, EventDataFilter.DataWithMetadata); 
}
private static void RegisterCacheLevelEvents() throws Exception {
	String cacheName = "demoCache";

	// Initialize an instance of the cache to begin performing operations:
	Cache cache = CacheManager.getCache(cacheName);

	CacheEventDescriptor eventDescriptor = cache.getMessagingService().addCacheNotificationListener(new CacheDataModificationListenerImpl(), EnumSet.allOf(EventType.class), EventDataFilter.DataWithMetadata);
}
 

Perché NCache?

Per le applicazioni .NET e Java, NCache è la scelta ideale per eseguire messaggi Pub/Sub con transazioni elevate, query continue ed eventi. Ecco perché.

  • Estremamente veloce e scalabile: NCache è estremamente veloce perché è totalmente in memoria e offre tempi di risposta inferiori al millisecondo. Inoltre, è una cache distribuita che ti consente di aggiungere server al cluster di cache ottenere una scalabilità lineare e gestire carichi di transazione estremiA differenza delle code di messaggistica basate su disco, NCache fornisce messaggistica Pub/Sub in memoria, garantendo una latenza inferiore al millisecondo per scenari ad alta produttività.
  • Elevata disponibilità tramite la replica dei dati: NCache replica in modo intelligente i dati nella cache senza compromettere le prestazioni. Quindi, non perderai alcun dato (eventi e messaggi) anche se un server cache si interrompe.
  • App .NET e Java Si scambiano messaggi: NCache consente alle applicazioni .NET e Java di condividere una piattaforma comune di messaggistica e flussi e di scambiarsi messaggi serializzando i messaggi in JSON o inviando documenti JSON come messaggi.

Con questo, puoi usare NCache come piattaforma di messaggistica potente ma semplice. Leggi di più su NCache dai collegamenti sottostanti e scarica una versione di prova di 30 giorni completamente funzionante.

Cosa fare dopo?

Domande frequenti (FAQ)

NCache Fornisce persistenza in memoria con replica. Sebbene progettato principalmente per velocità inferiori al millisecondo, garantisce un'elevata disponibilità replicando i messaggi su altri nodi del cluster. Supporta anche le sottoscrizioni durevoli, garantendo che i messaggi vengano archiviati fino alla riconnessione degli abbonati attivi o al raggiungimento della data di scadenza.

Pub/Sub è un modello di messaggistica disaccoppiato in cui editori e abbonati non si conoscono, ideale per microservizi e app di chat. Gli eventi cache (ItemAdded, ItemUpdated) sono strettamente associati alle modifiche dei dati all'interno della cache, destinati alla sincronizzazione delle cache locali o all'aggiornamento dei dati dell'interfaccia utente.

Sì, NCache Agisce come broker di messaggi distribuito in memoria ad alte prestazioni per microservizi .NET e Java. Elimina la necessità di broker esterni come RabbitMQ in scenari ad alta produttività, gestendo gli eventi di orchestrazione direttamente all'interno del livello di caching.

Perché NCache Elabora i messaggi interamente in memoria e offre una latenza inferiore al millisecondo anche in condizioni di carico estreme. I benchmark mostrano una scalabilità lineare, in grado di gestire milioni di operazioni al secondo man mano che si aggiungono nodi al cluster.

© Copyright Alachisoft 2002 - . Tutti i diritti riservati. NCache è un marchio registrato di Diyatech Corp.