NEST Elasticsearch .NETクライアントを使用したドキュメントのインデキシング
はじめに
NEST Elasticsearch .NETクライアントを使用して、Elasticsearchにドキュメントをインデックスする方法はいくつかあります。
このブログ記事では、一度に1つのドキュメントをインデキシングするシンプルな方法から、BulkObservableヘルパーを使用したより高度な方法まで、いくつかの手法を紹介します。
単一のドキュメント
NEST内では、ドキュメントはPOCO(Plain Old CLR Object)としてモデル化されます。以下に例を示します。
public class Person
{
public int Id { get; set; }
public 文字列 FirstName { get; set; }
public 文字列 LastName { get; set; }
}
Elasticsearchの単一ドキュメントを表すこのオブジェクトのインスタンスは、いくつかの異なるメソッドを使用してインデックスを作成できます。次のインスタンスを例として使用しましょう。
var person = new Person
{
Id = 1,
FirstName = "Martijn",
LastName = "Laarman"
};
IndexDocument<T>メソッドとIndexDocumentAsync<T>メソッドは、デフォルトパラメータを使用して、型Tの単一ドキュメントをインデックスする簡単な方法を提供します。このメソッド呼び出しの結果を調べることで、インデキシング操作が成功したかどうかを判断できます。
// IIndexResponseオブジェクトを返す同期メソッド
var indexResponse = client.IndexDocument(person);
// 待機可能なTask<IIndexResponse>を返す非同期メソッド
var indexResponseAsync = await client.IndexDocumentAsync(person);
// 同期操作の結果を検査します
if (!indexResponse.IsValid)
{
// If the request isn't valid, we can take action here
}
IsValidプロパティを使用して、対応が機能的に有効かどうかを確認できます。これは、リクエストで何らかの問題が発生したかどうかを単一のポイントで確認するためのNEST抽象化です。
ドキュメントのインデキシング時に追加のパラメーターを設定する必要がある場合は、fluent構文またはオブジェクト初期化構文を使用できます。これにより、インデキシングプロセスをより詳細に制御できるようになります。以下の例では、「people」という名前のインデックスにドキュメントをインデックスします。
// fluent構文
var fluentIndexResponse = client.Index(person, i => i.Index("people"));
// オブジェクト初期化構文
var initializerIndexResponse = client.Index(new IndexRequest<Person>(person, "people"));
複数のドキュメントをインデキシングする単純なアプローチとして、ループを作成して反復ごとに単一のドキュメントをインデキシングする方法がありますが、これは非常に非効率的であり、大規模なドキュメントコレクションではうまくスケールしません。
複数のドキュメント
Bulk APIは、複数のドキュメントのインデキシングに使用できます。まず、インデキシングするドキュメントのコレクションを作成しましょう:
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
};
IndexManyメソッドおよびIndexManyAsyncメソッドを使用すると、それぞれ同期または非同期で複数のドキュメントをインデキシングできます。これらのメソッドはNESTクライアント固有のものであり、クライアントのBulkメソッドおよびBulk APIへの呼び出しをラップすることで、多数のドキュメントをインデキシングするための便利なショートカットを提供します。
これらのメソッドはすべてのドキュメントを単一のHTTPリクエストでインデックスするため、非常に大規模なドキュメントコレクションの場合は、コレクションを多数の小さなバッチに分割し、複数のBulk呼び出しを発行する必要があることに注意してください。これを行う必要がある場合は、この投稿の後半で説明するBulkAllObservable<T>ヘルパーの使用を検討してください。
// IBulkResponseを返す同期メソッド
var indexManyResponse = client.IndexMany(people);
if (indexManyResponse.Errors)
{
// 対応でエラーを確認可能
foreach (var itemWithError in indexManyResponse.ItemsWithErrors)
{
// エラーがある場合、列挙して確認可能
Console.WriteLine("ドキュメント {0} のインデックスに失敗しました: {1}",
itemWithError.Id, itemWithError.Error);
}
}
// あるいは、ドキュメントを非同期でインデックスすることも可能
var indexManyAsyncResponse = await client.IndexManyAsync(people);
多数のドキュメントのインデキシングに対してよりきめ細かな制御が必要な場合は、BulkメソッドおよびBulkAsyncメソッドを使用し、記述子を使用してバルク呼び出しをカスタマイズできます。
上記のIndexManyメソッドと同様に、ドキュメントは単一のHTTPリクエストで_bulkエンドポイントに送信されます。そのため、HTTPリクエストの全体サイズを考慮する必要があります。大量のドキュメントをインデキシングする場合は、BulkAllObservable<T>ヘルパーを使用することをお勧めします。
// エラーを検査できるIBulkResponseを返します
var bulkIndexResponse = client.Bulk(b => b
.Index("people")
.IndexMany(people)
);
// 非同期バージョン
var asyncBulkIndexResponse = await client.BulkAsync(b => b
.Index("people")
.IndexMany(people)
);
BulkAllObservable<T> ヘルパー
BulkAllObservable<T>ヘルパーを使用すると、再試行、バックオフ、バッチ処理のメカニズムを気にすることなく、ドキュメントコレクションをインデキシングするという全体的な目的に集中できます。
BulkAllメソッドおよびBlockingSubscribeExtensionsのWait()拡張メソッドを使用して、複数のドキュメントをインデキシングできます。このヘルパーは、インデキシングの失敗時に自動的に再試行/バックオフを行う機能や、単一のHTTPリクエストでインデキシングするドキュメント数を制御する機能を公開しています。
以下の例では、各リクエストが元のインプットからバッチ処理された1000件のドキュメントをインデックス化します。ドキュメント数が非常に多い場合、それぞれ1000件のドキュメントを含む多数のHTTPリクエストが発生する可能性があります(最後の1回は、合計数に応じて1000件未満になる場合があります)。
このヘルパーはIEnumerable<T>コレクションを遅延列挙するため、ページ分割されたデータベースレコードから具体化されたドキュメントなど、大量のドキュメントを簡単にインデックスできます。
var bulkAllObservable = client.BulkAll(people, b => b
.Index("people")
// 再試行までの待機時間
.BackOffTime("30s")
// 失敗が発生した場合の再試行回数
.BackOffRetries(2)
// バルク操作完了後にインデックスをリフレッシュする
.RefreshOnCompleted()
// 同時に実行するバルク要求の数
.MaxDegreeOfParallelism(Environment.ProcessorCount)
// バルク要求あたりのアイテム数
.Size(1000)
)
// インデキシングを実行し、最大15分間待機します。
// BulkAll呼び出しは非同期ですが、これはブロッキング操作です。
.Wait(TimeSpan.FromMinutes(15), next =>
{
// do something on each response e.g. write number of batches indexed to console
});
BulkAllObservable<T>ヘルパーは、多数の高度な特徴を公開しています。
- BufferToBulkを使用すると、バルクリクエストがサーバーに送信される前に、そのリクエスト内の個々のオペレーションをカスタマイズできます。
- RetryDocumentPredicateを使用すると、インデックス作成に失敗したドキュメントを再試行するかどうかをきめ細かく制御できます。
- DroppedDocumentCallback: 再試行してもドキュメントがインデックスされない場合に呼び出されるデリゲートです。
client.BulkAll(people, b => b
.BufferToBulk((descriptor, list) =>
{
// バルク操作をディスパッチする前に
// 個別の操作をカスタマイズします
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) =>
{
// ドキュメントをインデックスできない場合にこのデリゲートが呼び出されます
Console.WriteLine($"Unable to index: {item} {person}");
})
);
インジェストノード
Elasticsearchは取り込みリクエストをインジェストノードへ自動的に再ルーティングするため、ルーティング情報を指定または構成する必要はありません。ただし、大量のインジェストを行っており、専用のインジェストノードがある場合は、クラスター内での余分なホップを避けるため、インデックスリクエストをこれらのノードに直接送信するのが合理的です。
これを実現する最もシンプルな方法は、専用の「インデキシング」クライアントインスタンスを作成し、それをインデキシングリクエストに使用することです。
// インジェストノードのリスト
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);
複雑なクラスター構成では、検査接続プールとNode述語を併用して、取り込み機能を持つNodeをフィルタリングする方が簡単な場合があります。これにより、クライアントを再構成することなくクラスターをカスタマイズできます。
// クラスターNodeのリスト
var pool = new SniffingConnectionPool(new []
{
new Uri("http://node1:9200"),
new Uri("http://node2:9200"),
new Uri("http://node3:9200")
});
// 取り込み機能を持つNodeのみを選択する述語
var settings = new ConnectionSettings(pool).NodePredicate(n => n.IngestEnabled);
var indexingClient = new ElasticClient(settings);
インジェストパイプライン
Person型を修正して、追加情報を含めます:
public class Person
{
public int Id { get; set; }
public 文字列 FirstName { get; set; }
public 文字列 LastName { get; set; }
public 文字列 IpAddress { get; set; }
public GeoIp GeoIp { get; set; }
}
public class GeoIp
{
public 文字列 CityName { get; set; }
public 文字列 ContinentName { get; set; }
public 文字列 CountryIsoCode { get; set; }
public GeoLocation Location { get; set; }
public 文字列 RegionName { get; set; }
}
インデックスされる前に受信値を操作するインジェストパイプラインを作成できます。アプリケーションが常に姓を大文字にすることを想定しており、イニシャルを独自のフィールドにインデックスする必要があると仮定します。また、人間が読める形式の場所に変換したいIPアドレスもあります。
この要件は、カスタムマッピングと取り込みパイプラインを作成することで実現できます。その後、これ以上の変更を加えることなく、新しいPersonタイプを使用できるようになります。
まず、インデックスとカスタムマッピングを作成します:
client.CreateIndex("people", c => c
.マッピング(ms => ms
.Map<Person>(p => p
//型からマッピングを自動作成
.AutoMap()
//AutoMap()から推論されたマッピングを上書き
.プロパティ(props => props
//イニシャルを格納するための追加フィールドを作成
.Keyword(t => t.Name("initials"))
//フィールドをIPアドレス型としてマッピング
.Ip(t => t.Name(dv => dv.IpAddress))
//GeoIpをオブジェクトとしてマッピング
.Object<GeoIp>(t => t.Name(dv => dv.GeoIp))
)
)
)
);
次に、バージョン6.7から同梱されているingest-geoipプラグインを活用して、取り込みパイプラインを作成します。
client.PutPipeline("person-pipeline", p => p
.Processors(ps => ps
//姓を大文字に変換
.Uppercase<Person>(s => s
.Field(t => t.LastName)
)
//Painlessスクリプトを使用して新しいフィールドに値を設定
.Script(s => s
.Lang("painless")
.Source("ctx.initials = ctx.firstName.substring(0,1) + ctx.lastName.substring(0,1)")
)
//ingest-geoipプラグインを使用して、提供されたIPアドレスからGeoIpオブジェクトを強化
.GeoIp<Person>(s => s
.Field(i => i.IpAddress)
.TargetField(i => i.GeoIp)
)
)
);
では、この新しいインデックスと取り込みパイプラインを使用してPersonインスタンスをインデックス化してみましょう。
var person = new Person
{
Id = 1,
FirstName = "Martijn",
LastName = "Laarman",
IpAddress = "139.130.4.5"
};
// 作成したパイプラインを使用してドキュメントをインデックス化
var indexResponse = client.Index(person, p => p
.Index("people")
.Pipeline("person-pipeline")
);
検索を実行すると、エンリッチされた値を含むインデックス済みのドキュメントが表示されます。
{
"took": 5,
"timed_out": false,
"_シャード": {
"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
}
}
]
}
}
パイプラインを指定すると、インデキシング時にドキュメントのエンリッチによるオーバーヘッドが追加されます。上記の例では、大文字変換とPainlessスクリプトの実行がこれに該当します。
大規模なバルクリクエストの場合、例外を回避するためにデフォルトのインデキシングタイムアウトを増やすのが賢明です。
client.Bulk(b => b
.Index("people")
.パイプライン("person-パイプライン")
//Elasticsearchサーバー側のタイムアウトを増やします
.Timeout("5m")
.IndexMany<Person>(people)
.RequestConfiguration(rc => rc
// リクエストを中断する前に、クライアント側のHTTPリクエストタイムアウトを増やします
.RequestTimeout(TimeSpan.FromMinutes(5))
)
);
まとめ
このブログ記事では、単一ドキュメントのインデキシングという単純なケースから、取り込みパイプラインを使用した複数ドキュメントのバルクインデキシングまでを解説しました。
ご自身のクラスタで試してみるか、 Elastic Cloud 上の Elasticsearch Service の14日間無料トライアルに登録してください。問題が発生した場合やご質問がある場合は、 Discuss フォーラムまでお問い合わせください。
NEST Elasticsearch .NETクライアントを使用したインデキシングの完全なドキュメントについては、ドキュメントをご覧ください。