Comment faire

Indexation de documents avec le client .NET NEST Elasticsearch

Introduction

Il existe plusieurs façons d'indexer des documents dans Elasticsearch en utilisant le client .NET NEST Elasticsearch.

Cet article de blog présentera certaines des méthodes simples, de l'indexation d'un seul document à la fois aux méthodes plus avancées utilisant l'assistant BulkObservable.

Documents uniques

Au sein de NEST, un document est modélisé sous forme de POCO (plain old CLR object) ; un exemple est donné ci-dessous :

public class Person
{
public int Id { get; set; }
public chaîne FirstName { get; set; }
public chaîne LastName { get; set; }
}

Une instance de cet objet, qui représente un document unique dans Elasticsearch, peut ensuite être indexée à l'aide de plusieurs méthodes différentes. Utilisons l'instance suivante comme exemple :

var person = new Person
{
Id = 1,
FirstName = "Martijn",
LastName = "Laarman"
};

Les<T> méthodes IndexDocument et IndexDocumentAsync<T> offrent un moyen simple d’indexer un document unique de type T, en utilisant les paramètres par défaut. Le résultat de cet appel de méthode peut être examiné pour déterminer si l’opération d’indexation a réussi.

// méthode synchrone qui renvoie un objet IIndexResponse
var indexResponse = client.IndexDocument(person);
// méthode asynchrone qui renvoie une tâche<IIndexResponse> IIndexResponse pouvant être attendue
var indexResponseAsync = await client.IndexDocumentAsync(person);
// Examiner le résultat de l’opération synchrone
if (!indexResponse.IsValid)
{
// If the request isn't valid, we can take action here
}

La propriété IsValid peut être utilisée pour vérifier si une réponse est fonctionnellement valide ou non. Il s’agit d’une abstraction NEST permettant de disposer d’un point unique pour vérifier si un problème est survenu lors de la requête.

Si vous devez définir des paramètres supplémentaires lors de l’indexation d’un document, vous pouvez utiliser la syntaxe fluent ou l’initialiseur d’objet. Cela vous donnera un contrôle plus fin sur le processus d’indexation. Dans l’exemple ci-dessous, nous allons indexer le document dans un index nommé « people ».

// syntaxe fluent
var fluentIndexResponse = client.Index(person, i => i.Index("people"));
// syntaxe d’initialiseur d’objet
var initializerIndexResponse = client.Index(new IndexRequest<Person>(person, "people"));

Une approche naïve de l’indexation de plusieurs documents consisterait à créer simplement une boucle pour indexer un seul document à chaque itération ; cependant, il s’agit d’une approche très inefficace qui ne pourra pas scaler correctement pour de grandes collections de documents.

Documents multiples

L’API Bulk peut être utilisée pour l’indexation de plusieurs documents. Tout d’abord, créons une collection de documents à indexer :

var people = new []
{
new Person
{
Id = 1,
FirstName = "Martijn",
LastName = "Laarman"
},
new Person
{
Id = 2,
FirstName = "Stuart",
LastName = "Cam"
},
new Person
{
Id = 3,
FirstName = "Russ",
LastName = "Cam"
}
// snip
};

Plusieurs documents peuvent être indexés à l'aide des méthodes IndexMany et IndexManyAsync, respectivement de manière synchrone ou asynchrone. Ces méthodes sont spécifiques au client NEST et encapsulent les appels à la méthode Bulk et à l'API Bulk du client, offrant un raccourci pratique pour l'indexation de nombreux documents.

Notez que ces méthodes indexent tous les documents dans une seule requête HTTP. Par conséquent, pour de très grandes collections de documents, vous devez partitionner la collection en plusieurs lots plus petits et émettre plusieurs appels Bulk. Lorsque vous devez effectuer cette opération, envisagez plutôt d'utiliser l'assistant BulkAllObservable<T>, décrit plus loin dans l'article.

// méthode synchrone qui renvoie une IBulkRéponse
var indexManyRéponse = client.IndexMany(people);
if (indexManyRéponse.Erreurs)
{
// la réponse peut être inspectée pour détecter des erreurs
foreach (var itemWithErreur in indexManyRéponse.ItemsWithErreurs)
{
// s'il y a des erreurs, elles peuvent être énumérées et inspectées
Console.WriteLine("Failed to index document {0}: {1}",
itemWithErreur.Id, itemWithErreur.Erreur);
}
}
// alternativement, les documents peuvent être indexés de manière asynchrone
var indexManyAsyncRéponse = await client.IndexManyAsync(people);

Si vous avez besoin d'un contrôle plus précis sur l'indexation d'un grand nombre de documents, vous pouvez utiliser les méthodes Bulk et BulkAsync et utiliser les descripteurs pour personnaliser les appels groupés.

Comme pour les méthodes IndexMany ci-dessus, les documents sont envoyés au point de terminaison _bulk dans une seule requête HTTP. Cela signifie qu'il faudra tenir compte de la taille globale de la requête HTTP. Pour l'indexation d'un grand nombre de documents, vous souhaiterez probablement utiliser l'assistant BulkAllObservable<T>.

// renvoie une IBulkResponse qui peut être inspectée pour détecter les erreurs
var bulkIndexResponse = client.Bulk(b => b
.Index("people")
.IndexMany(people)
);
// version asynchrone
var asyncBulkIndexResponse = await client.BulkAsync(b => b
.Index("people")
.IndexMany(people)
);

Assistant BulkAllObservable<T>

L'utilisation de l'assistant BulkAllObservable<T> vous permet de vous concentrer sur l'objectif global d'indexation d'une collection de documents, sans avoir à vous soucier des mécanismes de nouvelle tentative, d'interruption ou de traitement par lots.

Plusieurs documents peuvent être indexés à l'aide de la méthode BulkAll et de la méthode d'extension BlockingSubscribeExtensions Wait(). Cet assistant expose des fonctionnalités permettant de réessayer / de temporiser automatiquement en cas d'échec de l'indexation, et de contrôler le nombre de documents indexés dans une seule requête HTTP.

Dans l'exemple suivant, chaque requête indexe 1 000 documents, traités par lots à partir de l'entrée d'origine. En cas de grand nombre de documents, cela pourrait entraîner de nombreuses requêtes HTTP, chacune contenant 1 000 documents (la dernière requête peut en contenir moins, selon le nombre total).

L'assistant énumère de manière paresseuse une collection IEnumerable<T>, vous permettant d'indexer facilement un grand nombre de documents, tels que des documents matérialisés à partir d'enregistrements de base de données paginés.

var bulkAllObservable = client.BulkAll(people, b => b
.Index("people")
// durée d'attente entre les tentatives
.BackOffTime("30s")
// nombre de tentatives effectuées en cas d'échec
.BackOffRetries(2)
// actualiser l'index une fois l'opération en masse terminée
.RefreshOnCompleted()
// nombre de requêtes en masse simultanées à effectuer
.MaxDegreeOfParallelism(Environment.ProcessorCount)
// nombre d'éléments par requête en masse
.Size(1000)
)
// Effectuer l'indexation, en attendant jusqu'à 15 minutes.
// Bien que les appels BulkAll soient asynchrones, il s'agit d'une opération bloquante
.Wait(TimeSpan.FromMinutes(15), next =>
{
// do something on each response e.g. write number of batches indexed to console
});

L’assistant BulkAllObservable<T> expose un certain nombre de fonctionnalités avancées.

  1. BufferToBulk permet de personnaliser les opérations individuelles au sein de la requête en vrac avant qu'elle ne soit envoyée au serveur.
  2. RetryDocumentPredicate permet un contrôle précis pour décider si un document qui n'a pas pu être indexé doit faire l'objet d'une nouvelle tentative.
  3. DroppedDocumentCallback : dans le cas où un document n'est pas indexé, même après plusieurs tentatives, ce délégué est appelé.
client.BulkAll(people, b => b
.BufferToBulk((descriptor, list) =>
{
// personnaliser les opérations individuelles dans la requête
// bulk avant son envoi
foreach (var item in list)
{
// index each document into either even-index or odd-index
descriptor.Index<Person>(bi => bi
.Index(item.Id % 2 == 0 ? "even-index" : "odd-index")
.Document(item)
);
}
})
.RetryDocumentPredicate((item, person) =>
{
// decide if a document should be retried in the event of a failure
return item.Error.Index == "even-index" && person.FirstName == "Martijn";
})
.DroppedDocumentCallback((item, person) =>
{
// si un document ne peut pas être indexé, ce délégué est appelé
Console.WriteLine($"Impossible d’indexer : {item} {person}");
})
);

Nœuds d'ingestion

Étant donné qu'Elasticsearch redirige automatiquement les requêtes d'ingestion vers les nœuds d'ingestion, vous n'avez pas besoin de spécifier ou de configurer d'informations de routage. Cependant, si vous effectuez une ingestion intensive et que vous disposez de nœuds d'ingestion dédiés, il est judicieux d'envoyer les requêtes d'indexation directement vers ces nœuds afin d'éviter des sauts supplémentaires dans le cluster.

Le moyen le plus simple d'y parvenir est de créer une instance de client dédiée à l'« indexation » et de l'utiliser pour les requêtes d'indexation.

// liste des nœuds d'ingestion
var pool = new StaticConnectionPool(new []
{
new Uri("http://ingestnode1:9200"),
new Uri("http://ingestnode2:9200"),
new Uri("http://ingestnode3:9200")
});
var settings = new ConnectionSettings(pool);
var indexingClient = new ElasticClient(settings);

Dans les configurations de cluster complexes, il peut être plus simple d'utiliser un pool de connexions de type sniffing avec un prédicat de Node pour filtrer les Node dotés de capacités d'ingestion. Cela vous permet de personnaliser le cluster sans avoir à reconfigurer le client.

// list of cluster Node
var pool = new SniffingConnectionPool(new []
{
new Uri("http://node1:9200"),
new Uri("http://node2:9200"),
new Uri("http://node3:9200")
});
// predicate to select only Node with ingestion capabilities
var settings = new ConnectionSettings(pool).NodePredicate(n => n.ingestionEnabled);
var indexingClient = new ElasticClient(settings);

Pipelines d'ingestion

Modifions notre type Person pour inclure des informations supplémentaires :

public class Person
{
public int Id { get; set; }
public chaîne FirstName { get; set; }
public chaîne LastName { get; set; }
public chaîne IpAddress { get; set; }
public GeoIp GeoIp { get; set; }
}
public class GeoIp
{
public chaîne CityName { get; set; }
public chaîne ContinentName { get; set; }
public chaîne CountryIsoCode { get; set; }
public GeoLocation Location { get; set; }
public chaîne RegionName { get; set; }
}

Nous pouvons créer un pipeline d'ingestion qui manipule les valeurs entrantes avant qu'elles ne soient indexées. Supposons que notre application attende toujours que les noms de famille soient en majuscules et que les initiales soient indexées dans leur propre champ. Nous avons également une adresse IP que nous souhaitons convertir en un emplacement lisible par l'homme.

Nous pourrions répondre à cette exigence en créant un mapping personnalisé et un pipeline d'ingestion. Le nouveau type Person peut ensuite être utilisé sans effectuer d'autres modifications.

Tout d’abord, nous allons créer l’index et le mapping personnalisé :

client.CreateIndex("people", c => c
.Mappings(ms => ms
.Map<Person>(p => p
//créer automatiquement le mapping à partir du type
.AutoMap()
//remplacer tout mapping déduit d’AutoMap()
.Properties(props => props
// créer un champ supplémentaire pour stocker les initiales
.Keyword(t => t.Name("initials"))
//mapper le champ en tant que type adresse IP
.Ip(t => t.Name(dv => dv.IpAddress))
// mapper GeoIp en tant qu’objet
.Object<GeoIp>(t => t.Name(dv => dv.GeoIp))
)
)
)
);

Nous allons ensuite créer un pipeline d'ingestion en tirant parti du plug-in ingest-geoip, désormais intégré à la version 6.7.

client.PutPipeline("person-pipeline", p => p
.Processors(ps => ps
// mettre le nom de famille en majuscules
.Uppercase<Person>(s => s
.Field(t => t.LastName)
)
// utiliser un script painless pour remplir le nouveau champ
.Script(s => s
.Lang("painless")
.Source("ctx.initials = ctx.firstName.substring(0,1) + ctx.lastName.substring(0,1)")
)
// utiliser le plug-in ingest-geoip pour enrichir l'objet GeoIp à partir de l'adresse IP fournie
.GeoIp<Person>(s => s
.Field(i => i.IpAddress)
.TargetField(i => i.GeoIp)
)
)
);

Maintenant, indexons une instance Person en utilisant ce nouvel index et ce pipeline d'ingestion.

var person = new Person
{
Id = 1,
FirstName = "Martijn",
LastName = "Laarman",
IpAddress = "139.130.4.5"
};
// indexer le document en utilisant le pipeline créé
var indexResponse = client.Index(person, p => p
.Index("people")
.Pipeline("person-pipeline")
);

La recherche affiche désormais le document indexé avec les valeurs enrichies.

{
"took": 5,
"timed_out": false,
"_shards": {
"total": 1,
"successful": 1,
"skipped": 0,
"failed": 0
},
"hits": {
"total": 1,
"max_score": 1,
"hits": [
{
"_index": "people",
"_type": "person",
"_id": "1",
"_score": 1,
"_source": {
"firstName": "Martijn",
"lastName": "LAARMAN",
"initials": "ML",
"geoIp": {
"continent_name": "Oceania",
"region_iso_code": "AU-NSW",
"city_name": "Sydney",
"country_iso_code": "AU",
"region_name": "New South Wales",
"location": {
"lon": 151.2167,
"lat": -33.7333
}
},
"ipAddress": "139.130.4.5",
"id": 1
}
}
]
}
}

Lorsqu'un pipeline est spécifié, il y aura une surcharge supplémentaire liée à l'enrichissement des documents lors de l'indexation ; dans l'exemple donné ci-dessus, l'exécution de la mise en majuscules et du script Painless.

Pour les requêtes en masse volumineuses, il peut être prudent d'augmenter le délai d'expiration par défaut de l'indexation afin d'éviter les exceptions.

client.Bulk(b => b
.Index("people")
.Pipeline("person-pipeline")
//augmente le délai d'expiration côté serveur Elasticsearch
.Timeout("5m")
.IndexMany<Person>(people)
.RequestConfiguration(rc => rc
// augmente le délai d'expiration de la requête HTTP côté client, avant l'abandon de la requête
.RequestTimeout(TimeSpan.FromMinutes(5))
)
);

En résumé

Dans cet article de blog, nous avons couvert le cas simple de l'indexation d'un document unique, jusqu'à l'indexation groupée de plusieurs documents à l'aide de pipelines d'ingestion.

Essayez-le dans votre propre cluster ou déployez un essai gratuit de 14 jours du service Elasticsearch sur Elastic Cloud. Et si vous rencontrez des problèmes ou avez des questions, rendez-vous sur les forums de discussion Discuss.

Pour la documentation complète sur l'indexation à l'aide du client .NET NEST Elasticsearch, veuillez consulter notre documentation.