스트리밍 벌크 API

스트리밍 벌크 API (Streaming Bulk API)

한 번에 많은 문서를 처리하되, 배치 크기를 미리 가늠하지 않고도 안정적으로 처리하고 싶으시죠? 스트리밍 벌크 연산이 요청을 스트리밍하고 결과를 스트리밍 응답으로 받아 여러 문서를 추가, 업데이트, 삭제할 수 있게 해줘요.

출처: 문서

본문

2.17.0에서 도입

⚠️ 이 기능은 실험적(experimental)이며 프로덕션 환경에서 사용을 권장하지 않아요. 기능 진행 상황에 대한 업데이트나 피드백은 관련 GitHub issue를 참고하세요.

스트리밍 벌크 연산은 요청을 스트리밍하고 결과를 스트리밍 응답으로 받아 여러 문서를 추가, 업데이트, 삭제할 수 있게 해줘요. 기존 Bulk API와 비교하면, 스트리밍 수집은 배치 크기를 추정할 필요가 없고(이 값은 주어진 시점의 클러스터 운영 상태에 영향을 받아요), 많은 클라이언트와 클러스터 사이에 자연스러운 백프레셔를 적용해요. 스트리밍은 클라이언트와 클러스터의 기능에 따라 HTTP/2 또는 HTTP/1.1(chunked transfer encoding 사용)로 작동해요.

기본 HTTP transport 메서드는 스트리밍을 지원하지 않아요. transport-reactor-netty4 HTTP transport 플러그인을 설치하고 기본 HTTP transport 계층으로 사용해야 해요. transport-reactor-netty4 플러그인과 스트리밍 벌크 API 모두 실험적이에요.

엔드포인트

POST _bulk/stream
POST {index}/_bulk/stream

경로에 인덱스를 지정하면 요청 본문 청크에 인덱스를 포함할 필요가 없어요.

OpenSearch는 _bulk/stream 경로에 PUT 요청도 받지만, POST 사용을 강력히 권장해요. PUT의 허용된 용도—주어진 경로의 단일 리소스 추가 또는 대체—는 스트리밍 벌크 요청에는 맞지 않아요.

쿼리 파라미터

다음 표는 사용 가능한 쿼리 파라미터예요. 모든 쿼리 파라미터는 선택적이에요.

Parameter Data type Description
pipeline String 문서 전처리를 위한 파이프라인 ID예요.
refresh Enum 인덱싱 연산 후 영향을 받은 샤드를 refresh할지 여부예요. 기본값은 false예요. true는 변경 사항이 검색 결과에 즉시 나타나지만 클러스터 성능을 저하시켜요. wait_for는 refresh를 기다려요. 요청 반환은 더 오래 걸리지만 클러스터 성능은 저하되지 않아요.
require_alias Boolean 모든 작업이 인덱스가 아닌 인덱스 alias를 대상으로 하도록 true로 설정해요. 기본값은 false예요.
routing String 요청을 지정된 샤드로 라우팅해요.
timeout Time 요청이 반환되기를 얼마나 기다릴지예요. 기본값은 1m이에요.
type String (Deprecated) 유형을 지정하지 않은 문서의 기본 문서 유형이에요. 기본값은 _doc예요. 이 파라미터는 무시하고 모든 인덱스에 _doc 유형을 사용할 것을 강력히 권장해요.
wait_for_active_shards String OpenSearch가 벌크 요청을 처리하기 전에 사용 가능해야 하는 활성 샤드 수를 지정해요. 기본값은 1(primary 샤드만)이에요. all 또는 양의 정수로 설정해요. 1보다 큰 값은 replica가 필요해요. 예를 들어 값을 3으로 지정하면 요청이 성공하려면 인덱스에 2개의 replica가 2개의 추가 노드에 분산되어 있어야 해요.
batch_interval Time 벌크 연산을 데이터 노드로 보내기 전에 배치로 얼마나 오랫동안 누적할지 지정해요.
batch_size Time 벌크 연산을 데이터 노드로 보내기 전에 배치로 몇 개까지 누적할지 지정해요. 기본값은 1이에요.

요청 본문 필드

스트리밍 벌크 API 요청 본문은 Bulk API 요청 본문과 완전히 호환돼요. 각 벌크 연산(create/index/update/delete)은 별도의 청크로 전송돼요.

예시 요청

curl -X POST "http://localhost:9200/_bulk/stream" -H "Transfer-Encoding: chunked" -H "Content-Type: application/json" -d'
{ "delete": { "_index": "movies", "_id": "tt2229499" } }
{ "index": { "_index": "movies", "_id": "tt1979320" } }
{ "title": "Rush", "year": 2013 }
{ "create": { "_index": "movies", "_id": "tt1392214" } }
{ "title": "Prisoners", "year": 2013 }
{ "update": { "_index": "movies", "_id": "tt0816711" } }
{ "doc" : { "title": "World War Z" } }
'

예시 응답

배치 설정에 따라 각 스트리밍 응답 청크는 하나 또는 (배치) 여러 벌크 연산의 결과를 보고할 수 있어요. 예를 들어 배치하지 않은(기본값) 위 요청에 대해 스트리밍 응답은 다음과 같이 나타날 수 있어요.

{"took": 11, "errors": false, "items": [ { "index": {"_index": "movies", "_id": "tt1979320", "_version": 1, "result": "created", "_shards": { "total": 2 "successful": 1, "failed": 0 }, "_seq_no": 1, "_primary_term": 1, "status": 201 } } ] }
{"took": 2, "errors": true, "items": [ { "create": { "_index": "movies", "_id": "tt1392214", "status": 409, "error": { "type": "version_conflict_engine_exception", "reason": "[tt1392214]: version conflict, document already exists (current version [1])", "index": "movies", "shard": "0", "index_uuid": "yhizhusbSWmP0G7OJnmcLg" } } } ] }
{"took": 4, "errors": true, "items": [ { "update": { "_index": "movies", "_id": "tt0816711", "status": 404, "error": { "type": "document_missing_exception", "reason": "[_doc][tt0816711]: document missing", "index": "movies", "shard": "0", "index_uuid": "yhizhusbSWmP0G7OJnmcLg" } } } ] }

더 알아보기 (Learn more)

  • 스트리밍을 시작하기 전에 transport-reactor-netty4 플러그인 설치 여부를 확인하세요.
  • 실험적 기능이므로 충분히 검증한 뒤 사용하세요.