Indexação de documentos com o cliente NEST Elasticsearch .NET
Introdução
Existem várias maneiras de indexar documentos no Elasticsearch usando o cliente NEST Elasticsearch .NET.
Este post do blog demonstrará alguns dos métodos simples, desde a indexação de um único documento por vez até métodos mais avançados usando o auxiliar BulkObservable.
Documentos únicos
Dentro do NEST, um documento é modelado como POCO (plain old CLR object), um exemplo é fornecido abaixo:
public class Person
{
public int Id { get; set; }
public string FirstName { get; set; }
public string LastName { get; set; }
}
Uma instância deste objeto, que representa um único documento no Elasticsearch, pode então ser indexada usando alguns métodos diferentes. Vamos usar a seguinte instância como exemplo:
var person = new Person
{
Id = 1,
FirstName = "Martijn",
LastName = "Laarman"
};
Os métodos IndexDocument<T> e IndexDocumentAsync<T> fornecem uma maneira simples de indexar um único documento do tipo T, usando parâmetros padrão. O resultado desta chamada de método pode ser inspecionado para determinar se a operação de indexação foi bem-sucedida.
// método síncrono que retorna um objeto IIndexResponse
var indexResponse = client.IndexDocument(person);
// método assíncrono que retorna uma Task<IIndexResponse> que pode ser aguardada
var indexResponseAsync = await client.IndexDocumentAsync(person);
// Inspecione o resultado da operação síncrona
if (!indexResponse.IsValid)
{
// If the request isn't valid, we can take action here
}
A propriedade IsValid pode ser usada para verificar se uma resposta é funcionalmente válida ou não. Esta é uma abstração do NEST para ter um ponto único para verificar se algo deu errado com a solicitação.
Se você precisar definir parâmetros adicionais ao fazer a indexação de um documento, você pode usar a sintaxe fluente ou de inicializador de objeto. Isso lhe dará um controle mais refinado sobre o processo de indexação. No exemplo abaixo, indexaremos o documento em um índice chamado "people".
// sintaxe fluente
var fluentIndexResponse = client.Index(person, i => i.Index("people"));
// sintaxe de inicializador de objeto
var initializerIndexResponse = client.Index(new IndexRequest<Person>(person, "people"));
Uma abordagem ingênua para a indexação de vários documentos seria simplesmente criar um loop para indexar um único documento em cada iteração; no entanto, essa é uma abordagem muito ineficiente que não será redimensionada bem para grandes coleções de documentos.
Vários documentos
A bulk API pode ser usada para a indexação de múltiplos documentos. Primeiro, vamos criar uma coleção de documentos para indexar:
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
};
Vários documentos podem ser indexados usando os métodos IndexMany e IndexManyAsync, de forma síncrona ou assíncrona, respectivamente. Esses métodos são específicos do cliente NEST e encapsulam chamadas para o método Bulk e a bulk API do cliente, fornecendo um atalho conveniente para indexação de muitos documentos.
Observe que esses métodos indexam todos os documentos em uma única solicitação HTTP, portanto, para coleções de documentos muito grandes, você precisa particionar a coleção em muitos lotes menores e emitir várias chamadas Bulk. Quando você precisar fazer isso, considere usar o auxiliar BulkAllObservable<T>, descrito mais adiante no post.
// método síncrono que retorna uma IBulkResposta
var indexManyResposta = client.IndexMany(people);
if (indexManyResposta.Erros)
{
// a resposta pode ser inspecionada em busca de erros
foreach (var itemWithErro in indexManyResposta.ItemsWithErros)
{
// se houver erros, eles podem ser enumerados e inspecionados
Console.WriteLine("Falha ao indexar o documento {0}: {1}",
itemWithErro.Id, itemWithErro.Erro);
}
}
// alternativamente, os documentos podem ser indexados de forma assíncrona
var indexManyAsyncResposta = await client.IndexManyAsync(people);
Se você precisar de um controle mais granular sobre a indexação de muitos documentos, você pode usar os métodos Bulk e BulkAsync e usar os descritores para personalizar as chamadas em massa.
Assim como nos métodos IndexMany acima, os documentos são enviados para o endpoint _bulk em uma única solicitação HTTP. Isso significa que será necessário considerar o tamanho total da solicitação HTTP. Para a indexação de um grande número de documentos, você provavelmente desejará usar o auxiliar BulkAllObservable<T>.
// retorna uma IBulkResponse que pode ser inspecionada em busca de erros
var bulkIndexResponse = client.Bulk(b => b
.Index("people")
.IndexMany(people)
);
// versão assíncrona
var asyncBulkIndexResponse = await client.BulkAsync(b => b
.Index("people")
.IndexMany(people)
);
Auxiliar BulkAllObservable<T>
Usar o helper BulkAllObservable<T> permite que você foque no objetivo geral de indexar uma coleção de documentos, sem ter que se preocupar com mecanismos de repetição, backoff ou indexação em lote.
Vários documentos podem ser indexados usando o método BulkAll e o método de extensão BlockingSubscribeExtensions Wait(). Esse helper expõe a funcionalidade para tentar novamente/fazer backoff automaticamente no caso de uma falha de indexação, e para controlar o número de documentos indexados em uma única solicitação HTTP.
No exemplo a seguir, cada solicitação indexa 1.000 documentos, processados em lote a partir da entrada original. No caso de um grande número de documentos, isso pode resultar em muitas solicitações HTTP, cada uma contendo 1.000 documentos (a última solicitação pode conter menos, dependendo do número total).
O auxiliar enumera de forma preguiçosa uma coleção IEnumerable<T>, permitindo que você indexe um grande número de documentos facilmente, como documentos materializados a partir de registros de banco de dados paginados.
var bulkAllObservable = client.BulkAll(people, b => b
.Index("people")
// quanto tempo esperar entre as tentativas
.BackOffTime("30s")
// quantas tentativas são feitas se ocorrer uma falha
.BackOffRetries(2)
// atualizar o índice assim que a operação em massa for concluída
.RefreshOnCompleted()
// quantas solicitações em massa simultâneas fazer
.MaxDegreeOfParallelism(Environment.ProcessorCount)
// número de itens por solicitação em massa
.Size(1000)
)
// Realizar a indexação, aguardando até 15 minutos.
// Embora as chamadas BulkAll sejam assíncronas, esta é uma operação de bloqueio
.Wait(TimeSpan.FromMinutes(15), next =>
{
// do something on each response e.g. write number of batches indexed to console
});
O helper BulkAllObservable<T> expõe vários recursos avançados.
- O BufferToBulk permite a personalização de operações individuais dentro da solicitação em massa antes que ela seja enviada ao servidor.
- RetryDocumentPredicate permite um controle refinado ao decidir se um documento que falhou ao ser indexado deve ser tentado novamente.
- DroppedDocumentCallback: no caso de um documento não ser indexado, mesmo após tentativas, este delegado é chamado.
client.BulkAll(people, b => b
.BufferToBulk((descriptor, list) =>
{
// personalizar as operações individuais no bulk
// request antes de ser despachado
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) =>
{
// se um documento não puder ser indexado, este delegado é chamado
Console.WriteLine($"Não foi possível indexar: {item} {person}");
})
);
Nós de ingestão
Como o Elasticsearch redirecionará automaticamente as solicitações de ingestão para nós de ingestão, você não precisa especificar ou configurar nenhuma informação de roteamento. No entanto, se você estiver realizando uma ingestão pesada e tiver nós de ingestão dedicados, faz sentido enviar solicitações de índice diretamente para esses Node, para evitar saltos extras no cluster.
A maneira mais simples de conseguir isso é criar uma instância de cliente de "indexação" dedicada e usá-la para solicitações de indexação.
// lista de nós de ingestão
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);
Em configurações de cluster complexas, pode ser mais fácil usar um pool de conexões de sniffing junto com um predicado de Node para filtrar os Node que possuem recursos de ingestão. Isso permite que você personalize o cluster e não precise reconfigurar o cliente.
// lista de nós do cluster
var pool = new SniffingConnectionPool(new []
{
new Uri("http://node1:9200"),
new Uri("http://node2:9200"),
new Uri("http://node3:9200")
});
// predicado para selecionar apenas nós com recursos de ingestão
var settings = new ConnectionSettings(pool).NodePredicate(n => n.IngestEnabled);
var indexingClient = new ElasticClient(settings);
Pipelines de ingestão
Vamos modificar nosso tipo Person para incluir algumas informações adicionais:
public class Person
{
public int Id { get; set; }
public string FirstName { get; set; }
public string LastName { get; set; }
public string IpAddress { get; set; }
public GeoIp GeoIp { get; set; }
}
public class GeoIp
{
public string CityName { get; set; }
public string ContinentName { get; set; }
public string CountryIsoCode { get; set; }
public GeoLocation Location { get; set; }
public string RegionName { get; set; }
}
Podemos criar um pipeline de ingestão que manipula os valores recebidos antes que sejam indexados. Vamos supor que nossa aplicação sempre espere que os sobrenomes estejam em letras maiúsculas e que as iniciais sejam indexadas em seu próprio campo. Também temos um endereço IP que gostaríamos de converter em uma localização legível por humanos.
Poderíamos atender a esse requisito criando um mapeamento personalizado e um pipeline de ingestão. O novo tipo Person pode então ser usado sem fazer outras alterações.
Primeiro, criaremos o índice e o mapeamento personalizado:
client.CreateIndex("people", c => c
.Mappings(ms => ms
.Map<Person>(p => p
//cria automaticamente o mapeamento a partir do tipo
.AutoMap()
//substitui quaisquer mapeamentos inferidos do AutoMap()
.Properties(props => props
// cria um campo adicional para armazenar as iniciais
.Keyword(t => t.Name("initials"))
//mapeia o campo como tipo de endereço IP
.Ip(t => t.Name(dv => dv.IpAddress))
// mapeia GeoIp como objeto
.Object<GeoIp>(t => t.Name(dv => dv.GeoIp))
)
)
)
);
Em seguida, criaremos um pipeline de ingestão, aproveitando o plugin ingest-geoip incluído, agora disponível na versão 6.7.
client.PutPipeline("person-pipeline", p => p
.Processors(ps => ps
//colocar o sobrenome em maiúsculas
.Uppercase<Person>(s => s
.Field(t => t.LastName)
)
// usar um script Painless para preencher o novo campo
.Script(s => s
.Lang("painless")
.Source("ctx.initials = ctx.firstName.substring(0,1) + ctx.lastName.substring(0,1)")
)
// usar o plugin ingest-geoip para enriquecer o objeto GeoIp a partir do endereço IP fornecido
.GeoIp<Person>(s => s
.Field(i => i.IpAddress)
.TargetField(i => i.GeoIp)
)
)
);
Agora, vamos indexar uma instância de Person usando este novo índice e pipeline de ingestão.
var person = new Person
{
Id = 1,
FirstName = "Martijn",
LastName = "Laarman",
IpAddress = "139.130.4.5"
};
// indexar o documento usando o pipeline criado
var indexResponse = client.Index(person, p => p
.Index("people")
.Pipeline("person-pipeline")
);
A pesquisa agora mostra o documento indexado com os valores enriquecidos.
{
"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
}
}
]
}
}
Quando um pipeline é especificado, haverá a sobrecarga adicional de enriquecimento de documentos durante a indexação; no exemplo dado acima, a execução da conversão para maiúsculas e do script Painless.
Para grandes solicitações em massa, pode ser prudente aumentar o tempo limite de indexação padrão para evitar exceções.
client.Bulk(b => b
.Index("people")
.Pipeline("person-pipeline")
//aumenta o tempo limite no lado do servidor Elasticsearch
.Timeout("5m")
.IndexMany<Person>(people)
.RequestConfiguration(rc => rc
// aumenta o tempo limite da solicitação HTTP no cliente, antes de abortar a solicitação
.RequestTimeout(TimeSpan.FromMinutes(5))
)
);
Em resumo
Neste post do blog, cobrimos desde o caso simples de indexação de um único documento até a indexação em massa de vários documentos com pipelines de ingestão.
Experimente no seu próprio cluster ou inicie uma avaliação gratuita de 14 dias do Serviço Elasticsearch no Elastic Cloud. E se você executar algum problema ou tiver dúvidas, entre em contato nos fóruns do Discuss.
Para a documentação completa de indexação usando o cliente NEST Elasticsearch .NET, consulte nossa documentação.