Opensearch 커넥터

Opensearch 커넥터 (Opensearch Connector)

이 커넥터는 Opensearch Index 에 문서 액션을 요청할 수 있는 sink 를 제공합니다.

출처: 문서

본문

이 커넥터는 Opensearch Index 에 문서 액션을 요청할 수 있는 sink 를 제공합니다. 이 커넥터를 사용하려면 프로젝트에 다음 의존성을 추가하세요:

Opensearch 버전 Maven 의존성
1.x Flink 2.3 버전용 커넥터는 (아직) 사용할 수 없습니다.
2.x Flink 2.3 버전용 커넥터는 (아직) 사용할 수 없습니다.

기본적으로 Apache Flink Opensearch 커넥터는 1.3.x 클라이언트 라이브러리를 사용합니다. 예를 들어 2.x(또는 곧 나올 3.x) 클라이언트로 전환할 수 있는데, 이들은 JDK-11 이상 을 요구한다는 점에 유의하세요:

<dependencyManagement>
    <dependencies>
        <dependency>
            <groupId>org.opensearch</groupId>
            <artifactId>opensearch</artifactId>
            <version>2.5.0</version>
        </dependency>
        <dependency>
            <groupId>org.opensearch.client</groupId>
            <artifactId>opensearch-rest-high-level-client</artifactId>
            <version>2.5.0</version>
        </dependency>
    </dependencies>
</dependencyManagement>

스트리밍 커넥터는 현재 바이너리 배포에 포함되어 있지 않습니다. 클러스터 실행을 위해 라이브러리와 함께 프로그램을 패키징하는 방법은 여기 를 참고하세요.

Opensearch 설치

Opensearch 클러스터 설정 지침은 여기 에서 찾을 수 있습니다.

Opensearch Sink

아래 예제는 sink 를 구성하고 만드는 방법을 보여줍니다:

Java:

import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.connector.opensearch.sink.Opensearch2SinkBuilder;
import org.apache.flink.streaming.api.datastream.DataStream;

import org.apache.http.HttpHost;
import org.opensearch.action.index.IndexRequest;
import org.opensearch.client.Requests;

import java.util.HashMap;
import java.util.Map;

DataStream<String> input = ...;

input.sinkTo(
    new OpensearchSinkBuilder<String>()
        .setBulkFlushMaxActions(1) // Instructs the sink to emit after every element, otherwise they would be buffered
        .setHosts(new HttpHost("127.0.0.1", 9200, "http"))
        .setEmitter(
        (element, context, indexer) ->
        indexer.add(createIndexRequest(element)))
        .build());

private static IndexRequest createIndexRequest(String element) {
    Map<String, Object> json = new HashMap<>();
    json.put("data", element);

    return Requests.indexRequest()
        .index("my-index")
        .id(element)
        .source(json);
}

Scala:

import org.apache.flink.api.connector.sink.SinkWriter
import org.apache.flink.connector.opensearch.sink.{OpensearchSinkBuilder, RequestIndexer}
import org.apache.flink.streaming.api.datastream.DataStream
import org.apache.http.HttpHost
import org.opensearch.action.index.IndexRequest
import org.opensearch.client.Requests

val input: DataStream[String] = ...

input.sinkTo(
  new OpensearchSinkBuilder[String]
    .setBulkFlushMaxActions(1) // Instructs the sink to emit after every element, otherwise they would be buffered
    .setHosts(new HttpHost("127.0.0.1", 9200, "http"))
    .setEmitter((element: String, context: SinkWriter.Context, indexer: RequestIndexer) =>
    indexer.add(createIndexRequest(element)))
    .build())

def createIndexRequest(element: (String)): IndexRequest = {

  val json = Map(
    "data" -> element.asInstanceOf[AnyRef]
  )

  Requests.indexRequest.index("my-index").source(mapAsJavaMap(json))
}

이 예제는 각 수신 요소에 대해 단일 index 요청만 수행하는 것을 보여줍니다. 일반적으로 OpensearchEmitter 는 서로 다른 유형의 요청(예: DeleteRequest, UpdateRequest 등)을 수행하는 데 사용할 수 있습니다.

내부적으로 Flink Opensearch Sink 의 각 병렬 인스턴스는 클러스터에 액션 요청을 보내는 데 BulkProcessor 를 사용합니다. 이는 요소를 버퍼링한 뒤 클러스터에 벌크로 보냅니다. BulkProcessor 는 벌크 요청을 한 번에 하나씩 실행합니다. 즉, 버퍼된 액션의 두 개의 동시 플러시가 진행되지 않습니다.

Opensearch Sink 와 내결함성

Flink 의 체크포인트가 활성화되면 Flink Opensearch Sink 는 Opensearch 클러스터에 액션 요청의 최소 한 번(at-least-once) 전달을 보장합니다. 이는 체크포인트 시점에 BulkProcessor 의 모든 대기 중인 액션 요청을 기다림으로써 이루어집니다. 이는 체크포인트가 트리거되기 전의 모든 요청이 Opensearch 에 의해 성공적으로 승인된 후에만 sink 로 보내진 더 많은 레코드를 처리할 수 있음을 효과적으로 보장합니다.

체크포인트와 내결함성에 대한 자세한 내용은 내결함성 문서 에 있습니다.

내결함성 있는 Opensearch Sink 를 사용하려면 실행 환경에서 토폴로지의 체크포인팅을 활성화해야 합니다:

Java:

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(5000); // checkpoint every 5000 msecs

Scala:

val env = StreamExecutionEnvironment.getExecutionEnvironment()
env.enableCheckpointing(5000) // checkpoint every 5000 msecs

중요: 체크포인팅은 기본적으로 활성화되지 않지만 기본 전달 보장은 AT_LEAST_ONCE 입니다. 이로 인해 sink 는 완료되거나 BulkProcessor 가 자동으로 플러시될 때까지 요청을 버퍼링합니다. 기본적으로 BulkProcessor1000 개의 액션이 추가된 후 플러시합니다. 더 자주 플러시하도록 프로세서를 구성하려면 BulkProcessor 구성 섹션 을 참고하세요.

결정적 ID와 upsert 메서드가 있는 UpdateRequests 를 사용하면 커넥터에 AT_LEAST_ONCE 전달이 구성된 경우 Opensearch 에서 exactly-once 의미론을 달성할 수 있습니다.

실패한 Opensearch 요청 처리

Opensearch 액션 요청은 일시적으로 포화된 노드 큐 용량, 인덱싱할 말라포밍된(malformed) 문서 등 다양한 이유로 실패할 수 있습니다. Flink Opensearch Sink 는 사용자가 백오프(backoff) 정책을 지정해 요청을 재시도할 수 있게 합니다.

아래는 예제입니다:

Java:

DataStream<String> input = ...;

input.sinkTo(
    new OpensearchSinkBuilder<String>()
        .setHosts(new HttpHost("127.0.0.1", 9200, "http"))
        .setEmitter(
        (element, context, indexer) ->
        indexer.add(createIndexRequest(element)))
        // This enables an exponential backoff retry mechanism, with a maximum of 5 retries and an initial delay of 1000 milliseconds
        .setBulkFlushBackoffStrategy(FlushBackoffType.EXPONENTIAL, 5, 1000)
        .build());

Scala:

val input: DataStream[String] = ...

input.sinkTo(
  new OpensearchSinkBuilder[String]
    .setHosts(new HttpHost("127.0.0.1", 9200, "http"))
    .setEmitter((element: String, context: SinkWriter.Context, indexer: RequestIndexer) =>
    indexer.add(createIndexRequest(element)))
    // This enables an exponential backoff retry mechanism, with a maximum of 5 retries and an initial delay of 1000 milliseconds
    .setBulkFlushBackoffStrategy(FlushBackoffType.EXPONENTIAL, 5, 1000)
    .build())

위 예제는 리소스 제약(예: 큐 용량 포화)으로 실패한 요청을 sink 가 다시 추가하도록 합니다. 말라포밍된 문서 같은 다른 모든 실패의 경우 sink 는 실패합니다. BulkFlushBackoffStrategy(또는 FlushBackoffType.NONE)가 구성되지 않으면 sink 는 모든 종류의 오류에 대해 실패합니다.

중요: 실패 시 요청을 내부 BulkProcessor 에 다시 추가하면 체크포인트가 더 길어질 수 있습니다. 체크포인팅 시 sink 가 다시 추가된 요청이 플러시되는 것도 기다려야 하기 때문입니다. 예를 들어 FlushBackoffType.EXPONENTIAL 을 사용할 때 체크포인트는 Opensearch 노드 큐에 모든 대기 중인 요청에 충분한 용량이 생길 때까지, 또는 최대 재시도 횟수에 도달할 때까지 기다려야 합니다.

내부 Bulk Processor 구성

내부 BulkProcessor 는 버퍼된 액션 요청이 플러시되는 방식에 대한 동작을 OpensearchSinkBuilder 의 다음 메서드를 사용해 추가로 구성할 수 있습니다:

  • setBulkFlushMaxActions(int numMaxActions): 플러시 전 버퍼링할 최대 액션 수.
  • setBulkFlushMaxSizeMb(int maxSizeMb): 플러시 전 버퍼링할 최대 데이터 크기(메가바이트).
  • setBulkFlushInterval(long intervalMillis): 버퍼된 액션의 수나 크기와 무관하게 플러시하는 간격.

일시적인 요청 오류가 재시도되는 방식을 구성하는 것도 지원됩니다:

  • setBulkFlushBackoffStrategy(FlushBackoffType flushBackoffType, int maxRetries, long delayMillis): 백오프 지연 유형(CONSTANT 또는 EXPONENTIAL), 시도할 백오프 재시도 횟수, 백오프 지연량. 상수 백오프의 경우 이는 각 재시도 사이의 지연입니다. 지수 백오프의 경우 이는 초기 기본 지연입니다.

Opensearch 에 대한 자세한 정보는 여기 에서 찾을 수 있습니다.

Opensearch 커넥터를 Uber-Jar 로 패키징하기

Flink 프로그램 실행을 위해 모든 의존성을 포함하는 소위 uber-jar(실행 가능한 jar)를 빌드하는 것이 권장됩니다(여기 참고).

또는 커넥터의 jar 파일을 Flink 의 lib/ 폴더에 넣어 시스템 전역, 즉 실행되는 모든 작업에 사용 가능하게 할 수 있습니다.

더 알아보기 (Learn more)