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가 자동으로 플러시될 때까지 요청을 버퍼링합니다. 기본적으로BulkProcessor는1000개의 액션이 추가된 후 플러시합니다. 더 자주 플러시하도록 프로세서를 구성하려면 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/ 폴더에 넣어 시스템 전역, 즉 실행되는 모든 작업에 사용 가능하게 할 수 있습니다.