Agent 스트림 실행 API

Agent 스트림 실행 API (Execute Agent Stream API) (gRPC)

3.8에서 도입되었어요. gRPC Execute Agent Stream API는 gRPC를 통해 프로토콜 버퍼를 사용해 에이전트(agent) 실행을 스트리밍하는 이진 인터페이스를 제공해요. 에이전트가 출력을 생성할 때 서버가 응답 청크를 클라이언트로 스트리밍하므로, 클라이언트는 실행이 끝나기 전에 출력 처리를 시작할 수 있어요. 대규모 언어 모델(LLM)에서 토큰 단위로 출력을 생성하는 대화형 에이전트에는 스트리밍을 사용하세요.

에이전트 실행은 REST 또는 gRPC로 스트리밍할 수 있어요. 두 전송 모두 같은 증분 생성 출력을 반환하므로 클라이언트에 가장 잘 맞는 것을 선택하세요.

  • REST 스트리밍은 HTTP 위의 서버 전송 이벤트(SSE)를 사용하는데, 브라우저, 표준 HTTP 클라이언트, cURL 같은 명령줄 도구가 직접 지원해요. 이 기능은 실험적이며 프로덕션 환경에서의 사용은 권장되지 않아요. 자세한 내용은 Execute Agent Stream API를 참고하세요.
  • gRPC 스트리밍은 HTTP/2 위의 프로토콜 버퍼를 사용해요. 이 전송은 더 낮은 직렬화 오버헤드와 더 작은 페이로드, HTTP/2 흐름 제어와 연결 멀티플렉싱을 통한 네이티브 서버 스트리밍 의미론, 그리고 gRPC가 지원하는 어떤 언어로든 클라이언트를 생성할 수 있는 강력한 타입의 스키마를 제공해요.

스트리밍 에이전트 실행은 다음 외부 호스팅 모델을 사용하는 에이전트에 대해 지원돼요.

  • OpenAI Chat Completion
  • Amazon Bedrock Converse Stream

출처: 문서

본문

사전 요구 사항 (Prerequisites)

gRPC Execute Agent Stream API를 사용하기 전에 다음 사전 요구 사항을 충족했는지 확인하세요.

  • 클러스터에서 gRPC 전송을 활성화해요. 자세한 내용은 Using gRPC APIs를 참고하세요.
  • 클라이언트 쪽에서 ML Commons protobuf를 확보해요. protobuf를 얻는 방법은 Using gRPC APIs를 참고하세요.
  • 지원되는 스트리밍 모델을 사용하는 에이전트를 등록해요. 에이전트와 커넥터 구성은 Execute Agent Stream API를 참고하세요.

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

gRPC Execute Agent Stream API는 MLService 서비스에 있어요.

MLService 안의 ExecuteAgentStream 메서드를 호출해서 스트리밍 에이전트 실행 요청을 제출할 수 있어요. 이 메서드는 MlExecuteAgentStreamRequest를 받아 PredictResponse 메시지의 스트림을 반환해요.

ExecuteAgentStream은 서버 스트리밍 원격 프로시저 호출(RPC)이에요. 클라이언트가 단일 요청을 보내면 서버가 일련의 응답 메시지를 반환해요. 마지막 메시지는 is_last를 true로 설정하고, 서버는 그 후 스트림을 닫아요.

요청 필드 (Request fields)

gRPC Execute Agent Stream API는 다음 요청 필드를 지원해요.

MlExecuteAgentStreamRequest 필드 (MlExecuteAgentStreamRequest fields)

MlExecuteAgentStreamRequest 메시지는 다음 필드를 받아요.

필드 Protobuf 타입 필수 설명
agent_id string 필수 실행할 에이전트의 ID예요.
ml_execute_agent_stream_request_body MLExecuteAgentStreamRequestBody 필수 실행 파라미터를 담고 있는 요청 페이로드예요.

MLExecuteAgentStreamRequestBody 필드 (MLExecuteAgentStreamRequestBody fields)

MLExecuteAgentStreamRequestBody 메시지는 다음 필드를 받아요.

필드 Protobuf 타입 필수 설명
parameters Parameters 필수 에이전트에 전달되는 입력 파라미터예요.

Parameters 필드 (Parameters fields)

에이전트 실행의 경우 Parameters 메시지는 다음 필드를 받아요.

필드 Protobuf 타입 설명
question string 에이전트에 보내는 입력 질문이에요.

응답 필드 (Response fields)

서버는 일련의 PredictResponse 메시지를 스트리밍해요. 각 메시지는 생성된 출력의 청크 하나를 담고 다음 필드를 제공해요.

필드 Protobuf 타입 설명
inference_results repeated InferenceResults 청크에 대한 추론 결과예요.
inference_results.output repeated Output 각 추론 결과에 대한 출력 객체예요.
inference_results.output.name string 출력 필드의 이름(보통 response)이에요.
inference_results.output.result string memory_id와 parent_interaction_id 필드의 값이에요.
inference_results.output.data_as_map DataAsMap 청크에 대한 응답 콘텐츠와 메타데이터예요.
inference_results.output.data_as_map.content string 청크의 텍스트 콘텐츠예요. 전체 응답을 재구성하려면 청크들에 걸쳐 content 값을 이어붙여요.
inference_results.output.data_as_map.is_last bool 이것이 스트림의 마지막 청크인지 여부예요. true이면 더 이상 메시지가 전송되지 않아요.

예제 요청 (Example request)

다음 예제는 gRPC 요청 메시지의 JSON 표현을 보여줘요. agent_id와 question을 에이전트 구성과 일치하는 값으로 바꾸세요.

{
  "agent_id": "your_agent_id",
  "ml_execute_agent_stream_request_body": {
    "parameters": {
      "question": "List indices in my cluster"
    }
  }
}

다음 예제는 대화형 에이전트의 실행을 스트리밍하는 Java gRPC 클라이언트를 보여줘요. 에이전트 ID와 질문을 에이전트 구성과 일치하는 값으로 바꾸세요.

import org.opensearch.protobufs.*;
import org.opensearch.protobufs.services.MLServiceGrpc;
import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;

import java.util.Iterator;

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

        // Create a gRPC stub for ML operations
        MLServiceGrpc.MLServiceBlockingStub mlStub = MLServiceGrpc.newBlockingStub(channel);

        // Build the request parameters for the agent
        Parameters parameters = Parameters.newBuilder()
            .setQuestion("List indices in my cluster")
            .build();

        // Create the streaming agent execution request
        MlExecuteAgentStreamRequest request = MlExecuteAgentStreamRequest.newBuilder()
            .setAgentId("your_agent_id")
            .setMlExecuteAgentStreamRequestBody(MLExecuteAgentStreamRequestBody.newBuilder()
                .setParameters(parameters)
                .build())
            .build();

        // Execute the request and read the streamed response
        try {
            Iterator<PredictResponse> responses = mlStub.executeAgentStream(request);
            while (responses.hasNext()) {
                PredictResponse response = responses.next();
                for (InferenceResults results : response.getInferenceResultsList()) {
                    for (Output output : results.getOutputList()) {
                        DataAsMap chunk = output.getDataAsMap();
                        System.out.print(chunk.getContent());
                        if (chunk.getIsLast()) {
                            System.out.println("\n[stream complete]");
                        }
                    }
                }
            }
        } catch (io.grpc.StatusRuntimeException e) {
            System.err.println("gRPC execute agent stream request failed with status: " + e.getStatus());
            System.err.println("Error message: " + e.getMessage());
        }

        channel.shutdown();
    }
}

예제 응답 (Example response)

서버는 일련의 PredictResponse 메시지를 반환해요. 각 메시지는 content 필드에 생성된 텍스트의 청크 하나를 담고, 마지막 메시지는 isLast를 true로 설정해요. 다음 예제는 스트리밍된 청크의 JSON 표현을 보여줘요.

{
  "inferenceResults": [
    {
      "output": [
        {
          "name": "memory_id",
          "result": "6CMnkJ8BLGHoqtB13ipp"
        },
        {
          "name": "parent_interaction_id",
          "result": "6SMnkJ8BLGHoqtB13irB"
        },
        {
          "name": "response",
          "dataAsMap": {
            "content": "Hello",
            "isLast": false
          }
        }
      ]
    }
  ]
}
  • Execute Agent Stream API — 스트리밍 에이전트 실행의 REST 버전
  • Using gRPC APIs — gRPC 전송 구성과 클라이언트 요구 사항

더 알아보기 (Learn more)