Bulk API

Bulk API (gRPC)

3.0에서 도입되었어요. gRPC Bulk API는 인덱싱, 업데이트, 삭제 같은 여러 문서 작업을 단일 호출로 수행할 수 있는, HTTP Bulk API의 효율적인 이진 인코딩 대안이에요. 이 서비스는 프로토콜 버퍼(protocol buffers)를 사용하며 파라미터와 구조 면에서 REST API를 그대로 따릅니다.

출처: 문서

본문

사전 요구 사항 (Prerequisite)

gRPC 요청을 제출하려면 클라이언트 쪽에 protobuf 세트가 있어야 해요. protobuf를 얻는 방법은 Using gRPC APIs를 참고하세요.

gRPC 서비스와 메서드 (gRPC service and method)

gRPC Document API는 DocumentService에 있어요.

DocumentService 안의 Bulk gRPC 메서드를 호출해서 bulk 요청을 제출할 수 있어요. 이 메서드는 BulkRequest를 받아 BulkResponse를 반환해요.

문서 형식 (Document format)

gRPC에서 문서는 바이트로 제공되고 반환되어야 해요. gRPC 요청에서 문서를 제공하려면 Base64 인코딩을 사용하세요.

예를 들어 일반 Bulk API 요청의 다음 문서를 생각해보세요.

"doc":  "{\"title\": \"Inception\", \"year\": 2010}"

gRPC Bulk API 요청에서는 같은 문서를 Base64 인코딩으로 제공해요.

"doc": "eyJ0aX...MTB9"

BulkRequest 필드 (BulkRequest fields)

BulkRequest 메시지는 gRPC bulk 작업의 최상위 컨테이너예요. 다음 필드를 받아요.

필드 Protobuf 타입 설명
bulk_request_body repeated BulkRequestBody bulk 작업 목록이에요. 각각은 작업 유형(index/create/update/delete) 중 하나를 포함해요. 필수예요.
index string bulk_request_body에서 덮어쓰지 않는 한 모든 작업의 기본 인덱스예요. BulkRequest에서 인덱스를 지정하면 BulkRequestBody에 포함할 필요가 없어요. 선택 사항이에요.
x_source SourceConfigParam 응답에서 전체 _source, _source 없음, 또는 _source의 특정 필드만 반환할지 제어해요. 선택 사항이에요.
x_source_excludes repeated string source에서 제외할 필드예요. 선택 사항이에요.
x_source_includes repeated string source에서 포함할 필드예요. 선택 사항이에요.
pipeline string 전처리 ingest pipeline ID예요. 선택 사항이에요.
refresh Refresh 인덱싱 후 샤드를 새로고침할지 여부예요. 선택 사항이에요.
require_alias bool true이면 작업이 별칭을 대상으로 해야 해요. 선택 사항이에요.
routing string 샤드 할당을 위한 라우팅 값이에요. 선택 사항이에요.
timeout string 제한 시간(1m 같은)이에요. 선택 사항이에요.
type (Deprecated) string 문서 유형(항상 _doc)이에요. 선택 사항이에요.
wait_for_active_shards WaitForActiveShards 기다릴 최소 활성 샤드 수예요. 선택 사항이에요.
global_params GlobalParams 요청의 전역 파라미터예요. 선택 사항이에요.

BulkRequestBody 필드 (BulkRequestBody fields)

BulkRequestBody 메시지는 BulkRequest 안의 단일 문서 수준 작업을 나타내요. 다음 필드를 받아요.

필드 Protobuf 타입 설명
operation_container OperationContainer 수행할 작업(index, create, update, 또는 delete)이에요. 필수예요.
update_action UpdateAction 추가적인 update 전용 옵션이에요. 선택 사항이에요.
object bytes create와 index 작업에 사용하는 전체 문서 콘텐츠예요. 선택 사항이에요.

OperationContainer 필드 (OperationContainer fields)

OperationContainer 메시지는 정확히 하나의 작업 유형을 포함해요. 다음 필드를 받아요.

필드 Protobuf 타입 설명
index IndexOperation 문서를 인덱싱해요. 이미 존재하면 문서를 교체해요.
create WriteOperation 새 문서를 만들어요. 문서가 이미 존재하면 실패해요.
update UpdateOperation 문서를 부분 업데이트하거나 upsert/script 옵션을 사용해요.
delete DeleteOperation ID로 문서를 삭제해요.

UpdateAction 필드 (UpdateAction fields)

UpdateAction 메시지는 update 작업에 대한 추가 옵션을 제공해요. 다음 필드를 받아요.

필드 Protobuf 타입 설명
detect_noop bool true이면 문서 콘텐츠가 안 변했을 때 업데이트를 건너뛰어요. 선택 사항이에요. 기본값은 true예요.
doc bytes update 작업을 위한 부분 또는 전체 문서 데이터예요. 선택 사항이에요.
doc_as_upsert bool true이면 대상 문서가 존재하지 않을 때 문서를 전체 upsert 문서로 취급해요. update 작업에서만 유효해요. 선택 사항이에요.
script Script 문서에 적용할 스크립트예요(update와 함께 사용). 선택 사항이에요.
scripted_upsert bool true이면 문서가 존재하는지와 무관하게 스크립트를 실행해요. 선택 사항이에요.
upsert bytes 대상이 존재하지 않을 때 사용할 전체 문서예요. script와 함께 사용돼요. 선택 사항이에요.
x_source SourceConfig 문서 소스를 가져오거나 필터링하는 방법을 제어해요. 선택 사항이에요.

Create

WriteOperation은 문서가 아직 존재하지 않을 때만 새 문서를 추가해요.

문서 자체는 BulkRequestBody 메시지의 object 필드에 제공해야 해요.

다음 선택 필드도 제공할 수 있어요.

필드 Protobuf 타입 설명
x_id string 문서 ID예요. 생략하면 자동 생성돼요. 선택 사항이에요.
x_index string 대상 인덱스예요. BulkRequest에서 전역으로 설정되지 않았다면 필수예요. 선택 사항이에요.
routing string 샤드 배치를 제어하는 사용자 지정 라우팅 값이에요. 선택 사항이에요.
pipeline string 전처리 ingest pipeline ID예요. 선택 사항이에요.
require_alias bool true이면 모든 작업이 인덱스가 아닌 인덱스 별칭을 대상으로 해야 해요. 기본값은 false예요. 선택 사항이에요.
예제 요청 (Example request)

다음 예제는 create 작업이 있는 bulk 요청을 보여줘요. movies 인덱스에 ID tt1375666인 문서를 만들어요. Base64 인코딩으로 제공된 문서 콘텐츠는 {"title": "Inception", "year": 2010}을 나타내요.

{
  "index": "movies",
  "bulk_request_body": [
    {
      "operation_container": {
        "create": {
          "x_index": "movies",
          "x_id": "tt1375666"
        }
      },
      "object": "eyJ0aX...MTB9"
    }
  ]
}

Delete

DeleteOperation은 ID로 문서를 제거해요. 다음 필드를 받아요.

필드 Protobuf 타입 설명
x_id string 삭제할 문서의 ID예요. 필수예요.
x_index string 대상 인덱스예요. BulkRequest에서 전역으로 설정되지 않았다면 필수예요. 선택 사항이에요.
routing string 샤드 배치를 제어하는 사용자 지정 라우팅 값이에요. 선택 사항이에요.
if_primary_term int64 동시성 제어에 사용돼요. 문서의 프라이머리 텀이 이 값과 일치할 때만 작업이 실행돼요. 선택 사항이에요.
if_seq_no int64 동시성 제어에 사용돼요. 문서의 시퀀스 번호가 이 값과 일치할 때만 작업이 실행돼요. 선택 사항이에요.
version int64 동시성 제어를 위한 명시적 문서 버전이에요. 선택 사항이에요.
version_type VersionType 버전 일치 동작을 제어해요. 선택 사항이에요.
예제 요청 (Example request)

다음 예제는 delete 작업이 있는 bulk 요청을 보여줘요. movies 인덱스에서 ID tt1392214인 문서를 삭제해요.

{
  "index": "movies",
  "bulk_request_body": [
    {
      "operation_container": {
        "delete": {
          "x_index": "movies",
          "x_id": "tt1392214"
        }
      }
    }
  ]
}

Index

IndexOperation은 문서를 만들거나 덮어써요. ID를 제공하지 않으면 생성돼요.

문서 자체는 BulkRequestBody 메시지의 object 필드에 제공돼요.

다음 선택 필드도 제공할 수 있어요.

필드 Protobuf 타입 설명
x_id string 문서 ID예요. 생략하면 자동 생성돼요. 선택 사항이에요.
x_index string 대상 인덱스예요. BulkRequest에서 전역으로 설정되지 않았을 때만 필수예요.
routing string 샤드 배치를 제어하는 사용자 지정 라우팅 값이에요. 선택 사항이에요.
if_primary_term int64 동시성 제어에 사용돼요. 문서의 프라이머리 텀이 이 값과 일치할 때만 작업이 실행돼요. 선택 사항이에요.
if_seq_no int64 동시성 제어에 사용돼요. 문서의 시퀀스 번호가 이 값과 일치할 때만 작업이 실행돼요. 선택 사항이에요.
op_type OpType 작업 유형이에요. 덮어쓰기 동작을 제어해요. 유효한 값은 index(기본값)와 create예요. 선택 사항이에요.
version int64 동시성 제어를 위한 명시적 문서 버전이에요. 선택 사항이에요.
version_type VersionType 버전 일치 동작을 제어해요. 선택 사항이에요.
pipeline string 전처리 ingest pipeline ID예요. 선택 사항이에요.
require_alias bool true이면 모든 작업이 인덱스가 아닌 인덱스 별칭을 대상으로 해야 해요. 기본값은 false예요. 선택 사항이에요.
예제 요청 (Example request)

다음 예제는 index 작업이 있는 bulk 요청을 보여줘요. ID tt0468569인 Base64 인코딩 문서를 movies 인덱스에 인덱싱해요.

{
  "index": "movies",
  "bulk_request_body": [
    {
      "operation_container": {
        "index": {
          "x_index": "movies",
          "x_id": "tt0468569"
        }
      },
      "object": "eyJ0aX...MDh9"
    }
  ]
}

Update

UpdateOperation은 부분 문서 업데이트를 수행해요.

업데이트 옵션은 BulkRequestBody 메시지의 update_action 필드에 제공돼요.

다음 표에 나열된 모든 UpdateOperation 필드는 x_id를 제외하고 선택 사항이에요.

필드 Protobuf 타입 설명
x_id string 업데이트할 문서의 ID예요. 필수예요.
x_index string 대상 인덱스예요. BulkRequest에서 전역으로 설정되지 않았다면 필수예요. 선택 사항이에요.
routing string 샤드 배치를 제어하는 사용자 지정 라우팅 값이에요. 선택 사항이에요.
if_primary_term int64 동시성 제어에 사용돼요. 문서의 프라이머리 텀이 이 값과 일치할 때만 작업이 실행돼요. 선택 사항이에요.
if_seq_no int64 동시성 제어에 사용돼요. 문서의 시퀀스 번호가 이 값과 일치할 때만 작업이 실행돼요. 선택 사항이에요.
require_alias bool true이면 모든 작업이 인덱스가 아닌 인덱스 별칭을 대상으로 해야 해요. 기본값은 false예요. 선택 사항이에요.
retry_on_conflict int32 버전 충돌이 발생할 때 작업을 재시도할 횟수예요. 선택 사항이에요.
예제 요청 (Example request)

다음 예제는 update 작업이 있는 bulk 요청을 보여줘요. movies 인덱스에서 ID tt1375666인 문서를 {"year": 2011}로 업데이트해요.

{
  "index": "movies",
  "bulk_request_body": [
    {
      "operation_container": {
        "update": {
          "x_index": "movies",
          "x_id": "tt1375666"
        }
      },
      "update_action": {
        "doc": "eyJ5ZW...xMX0=",
        "detect_noop": true
      }
    }
  ]
}

Upsert

upsert 작업은 문서가 이미 존재하면 업데이트하고, 그렇지 않으면 제공된 문서 콘텐츠로 새 문서를 만들어요.

문서를 upsert하려면 UpdateOperation을 제공하고 BulkRequestBody에서 doc_as_upsert를 true로 지정해요. upsert할 문서는 doc 필드에 제공해야 해요.

예제 요청 (Example request)

다음 예제는 upsert 작업이 있는 bulk 요청을 보여줘요. movies 인덱스에서 ID tt1375666인 문서의 year 필드를 {"year": 2012}로 업데이트해요.

{
  "index": "movies",
  "bulk_request_body": [
    {
      "operation_container": {
        "update": {
          "x_index": "movies",
          "x_id": "tt1375666"
        }
      },
      "update_action": {
        "doc": "eyJ5ZW...xMn0=",
        "doc_as_upsert": true
      }
    }
  ]
}

Script

저장된 스크립트나 인라인 스크립트를 실행해 문서를 수정해요.

스크립트를 지정하려면 UpdateOperation과 BulkRequestBody에 script 필드를 제공해요.

예제 요청 (Example request)

다음 예제는 script 작업이 있는 bulk 요청을 보여줘요. movies 인덱스에서 ID tt1375666인 문서의 year 필드를 1만큼 증가시켜요.

{
  "index": "movies",
  "bulk_request_body": [
    {
      "operation_container": {
        "update": {
          "x_index": "movies",
          "x_id": "tt1375666"
        }
      },
      "update_action": {
        "script": {
          "inline": {
            "source": "ctx._source.year += 1",
            "lang": {
              "builtin": "BUILTIN_SCRIPT_LANGUAGE_PAINLESS"
            }
          }
        }
      }
    }
  ]
}

응답 필드 (Response fields)

gRPC Bulk API는 다음 응답 필드를 제공해요.

BulkResponse 필드 (BulkResponse fields)

BulkResponse 메시지는 Bulk gRPC 메서드에서 직접 반환되며 bulk 작업의 요약과 항목별 결과를 제공해요. 다음 필드를 포함해요.

필드 Protobuf 타입 설명
errors bool bulk 요청의 어떤 작업이라도 실패했는지 여부예요. 어떤 작업이 실패하면 응답의 errors 필드가 true가 돼요. 더 자세한 정보는 개별 Item 작업을 반복해서 확인할 수 있어요.
items repeated Item 제출된 순서대로 bulk 요청의 모든 작업 결과예요.
took int64 bulk 요청을 처리하는 데 걸린 시간(밀리초)이에요.
ingest_took int64 ingest pipeline을 통해 문서를 처리하는 데 걸린 시간(밀리초)이에요.

Item 필드 (Item fields)

응답의 각 Item은 요청의 단일 작업에 해당해요. 각 작업에 대해 다음 필드 중 하나만 제공돼요.

필드 Protobuf 타입 설명
create ResponseItem CreateOperation의 결과예요.
delete ResponseItem DeleteOperation의 결과예요.
index ResponseItem IndexOperation의 결과예요.
update ResponseItem UpdateOperation의 결과예요.

ResponseItem 필드 (ResponseItem fields)

각 ResponseItem은 요청의 단일 작업에 해당해요. 다음 필드를 포함해요.

필드 Protobuf 타입 설명
type string 문서 유형이에요.
id string 작업과 연결된 문서 ID예요.
index string 작업과 연결된 인덱스의 이름이에요. 데이터 스트림을 대상으로 했다면 이 값은 backing index예요.
status int32 작업에 대해 반환된 HTTP 상태 코드예요. (참고: 이 필드는 향후 gRPC 코드로 대체될 수 있어요.)
error ErrorCause 실패한 작업에 대한 추가 정보를 포함해요.
primary_term int64 문서에 할당된 프라이머리 텀이에요.
result string 작업 결과예요. 유효한 값은 created, deleted, updated예요.
seq_no int64 버전 순서를 유지하기 위해 문서에 할당된 시퀀스 번호예요.
shards ShardInfo 작업에 대한 샤드 정보예요(성공한 작업에만 반환).
version int64 문서 버전이에요(성공한 작업에만 반환).
forced_refresh bool true이면 작업 직후 문서가 즉시 보이도록 강제해요.
get InlineGetDictUserDefined 요청했다면 인라인 get에서 반환된 문서 소스를 포함해요.

InlineGetDictUserDefined 필드 (InlineGetDictUserDefined fields)

InlineGetDictUserDefined 메시지는 인라인 get 작업에서 반환된 문서 소스를 포함해요.

필드 Protobuf 타입 설명
metadata_fields optional ObjectMap 문서의 메타데이터 필드예요.
fields optional ObjectMap 문서의 저장된 필드예요.
found bool 문서가 존재하는지 여부예요.
x_seq_no optional int64 문서의 시퀀스 번호예요.
x_primary_term optional int64 문서의 프라이머리 텀이에요.
x_routing optional string 문서의 라우팅 값이에요.
x_source optional bytes 문서의 소스 데이터예요.

예제 응답 (Example response)

{
  "errors": false,
  "items": [
    {
      "index": {
        "x_id": "2",
        "x_index": "my_index",
        "status": 201,
        "x_primary_term": 1,
        "result": "created",
        "x_seq_no": 0,
        "x_shards": {
          "successful": 1,
          "total": 2
        },
        "x_version": 1,
        "forced_refresh": true
      }
    },
    {
      "create": {
        "x_id": "1",
        "x_index": "my_index",
        "status": 201,
        "x_primary_term": 1,
        "result": "created",
        "x_seq_no": 0,
        "x_shards": {
          "successful": 1,
          "total": 2
        },
        "x_version": 1,
        "forced_refresh": true
      }
    },
    {
      "update": {
        "x_id": "2",
        "x_index": "my_index",
        "status": 200,
        "x_primary_term": 1,
        "result": "updated",
        "x_seq_no": 1,
        "x_shards": {
          "successful": 1,
          "total": 2
        },
        "x_version": 2,
        "forced_refresh": true,
        "get": {
          "found": true,
          "x_seq_no": 1,
          "x_primary_term": 1,
          "x_source": "e30="
        }
      }
    },
    {
      "delete": {
        "x_id": "2",
        "x_index": "my_index",
        "status": 200,
        "x_primary_term": 1,
        "result": "deleted",
        "x_seq_no": 2,
        "x_shards": {
          "successful": 1,
          "total": 2
        },
        "x_version": 3,
        "forced_refresh": true
      }
    }
  ],
  "took": 87,
  "ingest_took": 0
}

Java gRPC 클라이언트 예제 (Java gRPC client example)

다음 예제는 예제 bulk gRPC 요청을 제출하고 bulk 응답에 오류가 있는지 확인하는 Java 클라이언트 프로그램을 보여줘요.

import org.opensearch.protobufs.*;
import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import com.google.protobuf.ByteString;

public class BulkClient {
    public static void main(String[] args) {
        ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 9400)
                .usePlaintext()
                .build();

        DocumentServiceGrpc.DocumentServiceBlockingStub stub = DocumentServiceGrpc.newBlockingStub(channel);

        // Create an index operation
        IndexOperation indexOp = IndexOperation.newBuilder()
                .setXIndex("my-index")
                .setXId("1")
                .build();

        BulkRequestBody indexBody = BulkRequestBody.newBuilder()
                .setOperationContainer(OperationContainer.newBuilder().setIndex(indexOp).build())
                .setObject(ByteString.copyFromUtf8("{\"field\": \"value\"}"))
                .build();

        // Create a delete operation
        DeleteOperation deleteOp = DeleteOperation.newBuilder()
                .setXIndex("my-index")
                .setXId("2")
                .build();

        BulkRequestBody deleteBody = BulkRequestBody.newBuilder()
                .setOperationContainer(OperationContainer.newBuilder().setDelete(deleteOp).build())
                .build();

        // Build the bulk request
        BulkRequest request = BulkRequest.newBuilder()
                .setIndex("my-index")
                .addBulkRequestBody(indexBody)
                .addBulkRequestBody(deleteBody)
                .build();

        // Execute the bulk request
        try {
            BulkResponse response = stub.bulk(request);

            // Handle the response
            System.out.println("Bulk errors: " + response.getErrors());
            System.out.println("Bulk took: " + response.getTook() + " ms");
            if (response.hasIngestTook()) {
                System.out.println("Ingest took: " + response.getIngestTook() + " ms");
            }

            // Process individual items
            for (Item item : response.getItemsList()) {
                if (item.hasIndex()) {
                    System.out.println("Index operation: " + item.getIndex().getStatus());
                } else if (item.hasDelete()) {
                    System.out.println("Delete operation: " + item.getDelete().getStatus());
                } else if (item.hasCreate()) {
                    System.out.println("Create operation: " + item.getCreate().getStatus());
                } else if (item.hasUpdate()) {
                    System.out.println("Update operation: " + item.getUpdate().getStatus());
                }
            }
        } catch (io.grpc.StatusRuntimeException e) {
            System.err.println("gRPC request failed with status: " + e.getStatus());
            System.err.println("Error message: " + e.getMessage());
        }

        channel.shutdown();
    }
}

Python gRPC 클라이언트 예제 (Python gRPC client example)

다음 예제는 Python 클라이언트 애플리케이션으로 같은 요청을 보내는 방법을 보여줘요.

먼저 pip로 opensearch-protobufs 패키지를 설치해요.

pip install opensearch-protobufs==1.2.0

다음 코드로 요청을 보내요.

import grpc

from opensearch.protobufs.schemas import *
from opensearch.protobufs.services import DocumentServiceStub

channel = grpc.insecure_channel(
    target="localhost:9400",
)

document_stub = DocumentServiceStub(channel)

# Add documents to a request body
requestBody = BulkRequestBody(
    operation_container=OperationContainer(index=IndexOperation())
)
requestBody.object = "{\"field\": \"value\"}".encode('utf-8')

# Append to a bulk request
request = BulkRequest()
request.index = "my-index"
request.bulk_request_body.append(requestBody)

# Send request and handle response
try:
    response = document_stub.Bulk(request=request)
    if response.items:
        print("Received {} response items".format(len(response.items)))
        print(response.items)
except grpc.RpcError as e:
    if e.code() == StatusCode.UNAVAILABLE:
        print("Failed to reach server: {}".format(e))
    elif e.code() == StatusCode.PERMISSION_DENIED:
        print("Permission denied: {}".format(e))
    elif e.code() == StatusCode.INVALID_ARGUMENT:
        print("Invalid argument: {}".format(e))
finally:
    channel.close()

더 알아보기 (Learn more)