Elasticsearch 커넥터
Elasticsearch 커넥터 (Elasticsearch Connector)
이 커넥터는 Elasticsearch 인덱스에 문서 액션을 요청할 수 있는 싱크(sink)를 제공합니다. Elasticsearch 설치 버전에 따라 프로젝트에 해당 의존성을 추가해야 합니다.
출처: 문서
본문
이 커넥터는 Elasticsearch 인덱스에 문서 액션을 요청할 수 있는 싱크를 제공합니다. 이 커넥터를 사용하려면 Elasticsearch 설치 버전에 따라 프로젝트에 다음 의존성 중 하나를 추가하세요:
| Elasticsearch 버전 | Maven 의존성 |
|---|---|
| 6.x | Flink 버전 2.3용 커넥터는 아직 없습니다. |
| 7.x | Flink 버전 2.3용 커넥터는 아직 없습니다. |
| 8.x | Flink 버전 2.3용 커넥터는 아직 없습니다. |
PyFlink 작업에서 사용하려면 다음 의존성이 필요합니다:
| 버전 | PyFlink JAR |
|---|---|
| flink-connector-elasticsearch6 | Flink 버전 2.3용 SQL jar는 아직 없습니다. |
| flink-connector-elasticsearch7 | Flink 버전 2.3용 SQL jar는 아직 없습니다. |
PyFlink에서 JAR을 사용하는 방법에 대한 자세한 내용은 Python dependency management를 참고하세요.
참고로 스트리밍 커넥터는 현재 바이너리 배포판의 일부가 아닙니다. 클러스터 실행을 위해 프로그램을 라이브러리와 함께 패키징하는 방법은 여기에서 확인하세요.
Elasticsearch 설치 (Installing Elasticsearch)
Elasticsearch 클러스터 설정에 대한 지침은 여기에서 찾을 수 있습니다.
Elasticsearch 싱크 (Elasticsearch Sink)
아래 예시는 싱크를 구성하고 생성하는 방법을 보여줍니다.
Elasticsearch 6 (Java):
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.connector.elasticsearch.sink.Elasticsearch6SinkBuilder;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.http.HttpHost;
import org.elasticsearch.action.index.IndexRequest;
import org.elasticsearch.client.Requests;
import java.util.HashMap;
import java.util.Map;
DataStream<String> input = ...;
input.sinkTo(
new Elasticsearch6SinkBuilder<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")
.type("my-type")
.id(element)
.source(json);
}
Elasticsearch 7 (Java):
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.connector.elasticsearch.sink.Elasticsearch7SinkBuilder;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.http.HttpHost;
import org.elasticsearch.action.index.IndexRequest;
import org.elasticsearch.client.Requests;
import java.util.HashMap;
import java.util.Map;
DataStream<String> input = ...;
input.sinkTo(
new Elasticsearch7SinkBuilder<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);
}
Elasticsearch 8 (Java):
import co.elastic.clients.elasticsearch.core.bulk.IndexOperation;
import org.apache.flink.connector.elasticsearch.sink.Elasticsearch8AsyncSinkBuilder;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.http.HttpHost;
import java.util.HashMap;
import java.util.Map;
DataStream<String> input = ...;
input.sinkTo(
Elasticsearch8AsyncSinkBuilder.<String>builder()
.setHosts(new HttpHost("127.0.0.1", 9200, "http"))
.setMaxBatchSize(1) // Instructs the sink to emit after every element, otherwise they would be buffered
.setElementConverter((element, ctx) -> {
Map<String, Object> json = new HashMap<>();
json.put("data", element);
return new IndexOperation.Builder<>()
.id(element)
.document(json)
.index("my-index")
.build();
})
.build());
Elasticsearch 6 (Scala):
import org.apache.flink.api.connector.sink.SinkWriter
import org.apache.flink.connector.elasticsearch.sink.{Elasticsearch6SinkBuilder, RequestIndexer}
import org.apache.flink.streaming.api.datastream.DataStream
import org.apache.http.HttpHost
import org.elasticsearch.action.index.IndexRequest
import org.elasticsearch.client.Requests
import scala.collection.JavaConverters._
val input: DataStream[String] = ...
input.sinkTo(
new Elasticsearch6SinkBuilder[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").`type`("my-type").source(mapAsJavaMap(json))
}
Elasticsearch 7 (Scala):
import org.apache.flink.api.connector.sink.SinkWriter
import org.apache.flink.connector.elasticsearch.sink.{Elasticsearch7SinkBuilder, RequestIndexer}
import org.apache.flink.streaming.api.datastream.DataStream
import org.apache.http.HttpHost
import org.elasticsearch.action.index.IndexRequest
import org.elasticsearch.client.Requests
import scala.collection.JavaConverters._
val input: DataStream[String] = ...
input.sinkTo(
new Elasticsearch7SinkBuilder[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))
}
Elasticsearch 8 (Scala):
import co.elastic.clients.elasticsearch.core.bulk.IndexOperation
import org.apache.flink.connector.elasticsearch.sink.Elasticsearch8AsyncSinkBuilder
import org.apache.flink.streaming.api.datastream.DataStream
import org.apache.http.HttpHost
import org.apache.flink.api.connector.sink.SinkWriter
import scala.collection.JavaConverters.mapAsJavaMap
val input: DataStream[String] = ...
input.sinkTo(
Elasticsearch8AsyncSinkBuilder.builder[String]()
.setHosts(new HttpHost("127.0.0.1", 9200, "http"))
.setMaxBatchSize(1) // Instructs the sink to emit after every element, otherwise they would be buffered
.setElementConverter { (element, ctx) =>
val json = Map("data" -> element.asInstanceOf[AnyRef])
new IndexOperation.Builder[Object]()
.id(element)
.document(mapAsJavaMap(json))
.index("my-index")
.build()
}
.build())
Elasticsearch 6 정적 인덱스 (Python):
from pyflink.datastream.connectors.elasticsearch import Elasticsearch6SinkBuilder, ElasticsearchEmitter
env = StreamExecutionEnvironment.get_execution_environment()
env.add_jars(ELASTICSEARCH_SQL_CONNECTOR_PATH)
input = ...
# The set_bulk_flush_max_actions instructs the sink to emit after every element, otherwise they would be buffered
es6_sink = Elasticsearch6SinkBuilder() \
.set_bulk_flush_max_actions(1) \
.set_emitter(ElasticsearchEmitter.static_index('foo', 'id', 'bar')) \
.set_hosts(['localhost:9200']) \
.build()
input.sink_to(es6_sink).name('es6 sink')
Elasticsearch 6 동적 인덱스 (Python):
from pyflink.datastream.connectors.elasticsearch import Elasticsearch6SinkBuilder, ElasticsearchEmitter
env = StreamExecutionEnvironment.get_execution_environment()
env.add_jars(ELASTICSEARCH_SQL_CONNECTOR_PATH)
input = ...
es_sink = Elasticsearch6SinkBuilder() \
.set_emitter(ElasticsearchEmitter.dynamic_index('name', 'id', 'bar')) \
.set_hosts(['localhost:9200']) \
.build()
input.sink_to(es6_sink).name('es6 dynamic index sink')
Elasticsearch 7 정적 인덱스 (Python):
from pyflink.datastream.connectors.elasticsearch import Elasticsearch7SinkBuilder, ElasticsearchEmitter
env = StreamExecutionEnvironment.get_execution_environment()
env.add_jars(ELASTICSEARCH_SQL_CONNECTOR_PATH)
input = ...
# The set_bulk_flush_max_actions instructs the sink to emit after every element, otherwise they would be buffered
es7_sink = Elasticsearch7SinkBuilder() \
.set_bulk_flush_max_actions(1) \
.set_emitter(ElasticsearchEmitter.static_index('foo', 'id')) \
.set_hosts(['localhost:9200']) \
.build()
input.sink_to(es7_sink).name('es7 sink')
Elasticsearch 7 동적 인덱스 (Python):
from pyflink.datastream.connectors.elasticsearch import Elasticsearch7SinkBuilder, ElasticsearchEmitter
env = StreamExecutionEnvironment.get_execution_environment()
env.add_jars(ELASTICSEARCH_SQL_CONNECTOR_PATH)
input = ...
es7_sink = Elasticsearch7SinkBuilder() \
.set_emitter(ElasticsearchEmitter.dynamic_index('name', 'id')) \
.set_hosts(['localhost:9200']) \
.build()
input.sink_to(es7_sink).name('es7 dynamic index sink')
예시는 들어오는 각 요소에 대해 단일 인덱스 요청만 수행하는 것을 보여줍니다. 일반적으로 ElasticsearchEmitter는 다른 유형의 요청(예: DeleteRequest, UpdateRequest 등)을 수행하는 데 사용할 수 있습니다.
내부적으로 Flink Elasticsearch Sink의 각 병렬 인스턴스는 BulkProcessor를 사용하여 클러스터에 액션 요청을 보냅니다. 이는 클러스터에 일괄 전송하기 전에 요소를 버퍼링합니다. BulkProcessor는 한 번에 하나씩 벌크 요청을 실행합니다. 즉, 버퍼링된 액션의 동시 플러시가 진행되지 않습니다.
Elasticsearch 싱크와 내결함성 (Elasticsearch Sinks and Fault Tolerance)
Flink의 체크포인팅이 활성화된 상태에서 Flink Elasticsearch Sink는 액션 요청의 at-least-once 전달을 Elasticsearch 클러스터에 보장합니다. 체크포인트 시점에 BulkProcessor의 모든 보류 중인 액션 요청을 기다림으로써 이를 수행합니다. 이는 체크포인트가 트리거되기 전의 모든 요청이 싱크로 보내진 더 많은 레코드를 처리하기 전에 Elasticsearch가 성공적으로 확인(acknowledge)했음을 효과적으로 보장합니다. 참고로 Elasticsearch 8 싱크는 Async Sink를 사용하여 at-least-once 의미론을 제공합니다. 자세한 내용은 FLIP-171 참고.
체크포인트와 내결함성에 대한 자세한 내용은 fault tolerance docs에 있습니다.
백오프 (Backoff)
input = ...
# This enables an exponential backoff retry mechanism, with a maximum of 5 retries and an initial delay of 1000 milliseconds
es7_sink = Elasticsearch7SinkBuilder() \
.set_bulk_flush_backoff_strategy(FlushBackoffType.EXPONENTIAL, 5, 1000) \
.set_emitter(ElasticsearchEmitter.static_index('foo', 'id')) \
.set_hosts(['localhost:9200']) \
.build()
input.sink_to(es7_sink).name('es7 sink')
위 예시는 리소스 제약(예: 큐 용량 포화)으로 인해 실패한 요청을 싱크가 다시 추가하도록 합니다. 형식이 잘못된 문서 같은 다른 모든 실패에 대해서는 싱크가 실패합니다. BulkFlushBackoffStrategy(또는 FlushBackoffType.NONE)가 구성되지 않으면 싱크는 어떤 종류의 오류에도 실패합니다.
중요: 실패 시 요청을 내부 BulkProcessor에 다시 추가하면 체크포인트가 더 길어질 수 있습니다. 체크포인트할 때 싱크가 재추가된 요청이 플러시되기를 기다려야 하기 때문입니다. 예를 들어 FlushBackoffType.EXPONENTIAL을 사용하면 체크포인트는 Elasticsearch 노드 큐에 모든 보류 요청에 충분한 용량이 생길 때까지, 또는 최대 재시도 횟수에 도달할 때까지 기다려야 합니다.
내부 Bulk Processor 구성 (Configuring the Internal Bulk Processor)
내부 BulkProcessor는 버퍼링된 액션 요청이 플러시되는 방식에 대한 동작을 Elasticsearch6SinkBuilder의 다음 메서드로 추가 구성할 수 있습니다:
- setBulkFlushMaxActions(int numMaxActions): 플러시 전 버퍼링할 최대 액션 수.
- setBulkFlushMaxSizeMb(int maxSizeMb): 플러시 전 버퍼링할 최대 데이터 크기(메가바이트).
- setBulkFlushInterval(long intervalMillis): 버퍼링된 액션의 양이나 크기와 무관하게 플러시하는 간격.
일시적 요청 오류가 재시도되는 방식도 지원됩니다:
- setBulkFlushBackoffStrategy(FlushBackoffType flushBackoffType, int maxRetries, long delayMillis):
CONSTANT또는EXPONENTIAL중 백오프 지연 유형, 시도할 백오프 재시도 횟수, 백오프 지연 양. 상수 백오프의 경우 각 재시도 사이의 지연입니다. 지수 백오프의 경우 초기 기준 지연입니다.
Elasticsearch에 대한 자세한 정보는 여기에서 찾을 수 있습니다.
내부 Writer 구성 (Configuring the Internal Writer)
참고로 Elasticsearch 8 싱크는 레거시 BulkProcessor 대신 AsyncSinkWriter를 사용합니다. Elasticsearch8AsyncSinkBuilder의 다음 메서드로 구성할 수 있습니다:
- setMaxBatchSize(int maxBatchSize): 배치의 최대 레코드 수.
- setMaxInFlightRequests(int maxInFlightRequests): 허용되는 in-flight 요청의 최대 수.
- setMaxBufferedRequests(int maxBufferedRequests): 싱크에 버퍼링될 수 있는 최대 레코드 수.
- setMaxBatchSizeInBytes(int maxBatchSizeInBytes): 배치가 될 수 있는 최대 크기(바이트).
- setMaxTimeInBufferMS(int maxTimeInBufferMS): 레코드가 플러시되기 전 싱크에 머무를 수 있는 최대 시간(밀리초).
- setMaxRecordSizeInBytes(int maxRecordSizeInBytes): 싱크가 수락할 최대 레코드 크기. 이보다 큰 레코드는 자동으로 거부됩니다.
Elasticsearch 커넥터를 Uber-Jar로 패키징 (Packaging the Elasticsearch Connector into an Uber-Jar)
Flink 프로그램 실행 시 모든 의존성을 포함하는 소위 uber-jar(실행 가능한 jar)를 빌드하는 것이 권장됩니다(여기 참고).
또는 커넥터의 jar 파일을 Flink의 lib/ 폴더에 넣어 시스템 전역(즉, 실행되는 모든 작업에 대해)으로 사용 가능하게 할 수 있습니다.
ElasticsearchEmitter는 동적 인덱스 이름과 Delete/Update 연산을 지원합니다(상세 예시는 커넥터의 테스트 소스 참고).