Elasticsearch Sink 커넥터
Elasticsearch Sink 커넥터
Elasticsearch sink 커넥터는 Pulsar 토픽에서 메시지를 가져와 Elasticsearch 인덱스에 저장하는 커넥터예요. Pulsar의 데이터를 검색·분석 엔진인 Elasticsearch에 색인할 때 사용해요.
참고: 모든 Pulsar 커넥터는 download page에서 내려받을 수 있어요.
출처: 문서
본문
요구 사항 (Requirements)
Elasticsearch sink 커넥터를 배포하려면 다음이 필요해요.
- Elasticsearch 7 (Elasticsearch 8은 향후 지원 예정)
- OpenSearch 1.x
기능 (Feature)
데이터 처리 (Handle data)
Pulsar 2.9.0부터 Elasticsearch sink 커넥터는 다음과 같은 동작 방식을 가져요. 그중 하나를 선택할 수 있어요.
| 이름 | 설명 |
|---|---|
| Raw processing | sink가 토픽에서 읽은 원본(raw) 콘텐츠를 그대로 Elasticsearch로 전달해요. 기본 동작이에요. Raw processing은 Pulsar 2.8.x에서도 이미 사용할 수 있었어요. |
| Schema aware | sink가 스키마를 사용하고 AVRO, JSON, KeyValue 스키마 타입을 처리하면서 콘텐츠를 Elasticsearch 도큐먼트에 매핑해요. schemaEnable을 true로 설정하면 sink가 메시지의 내용을 해석하고, Elasticsearch의 특수 _id 필드로 사용되는 기본 키를 정의할 수 있어요. |
여러 인덱스 매핑 (Map multiple indexes)
Pulsar 2.9.0부터 indexName 프로퍼티는 더 이상 필수가 아니에요. 생략하면 sink는 Pulsar 토픽 이름을 따른 인덱스 이름에 써요.
벌크 쓰기 활성화 (Enable bulk writes)
Pulsar 2.9.0부터 bulkEnabled 프로퍼티를 true로 설정해 벌크 쓰기를 사용할 수 있어요.
TLS로 보안 연결 활성화 (Enable secure connections via TLS)
Pulsar 2.9.0부터 TLS로 보안 연결을 활성화할 수 있어요.
구성 (Configuration)
Elasticsearch sink 커넥터의 구성에는 다음과 같은 프로퍼티가 있어요.
프로퍼티 (Property)
| 이름 | 타입 | 필수 | 기본값 | 설명 |
|---|---|---|---|---|
| elasticSearchUrl | String | true | " " (빈 문자열) | 커넥터가 연결하는 Elasticsearch 클러스터의 URL이에요. |
| indexName | String | false | " " (빈 문자열) | 커넥터가 메시지를 쓰는 인덱스 이름이에요. 기본값은 토픽 이름이에요. 이름에 날짜 형식을 허용해 %{+<date-format>} 패턴으로 이벤트 시간 기반 인덱스를 지원해요. 예를 들어 레코드의 이벤트 시간이 1645182000000L이고 indexName이 logs-%{+yyyy-MM-dd}라면, 포맷된 인덱스 이름은 logs-2022-02-18이 돼요. |
| schemaEnable | Boolean | false | false | Schema Aware 모드를 켤지 여부예요. |
| createIndexIfNeeded | Boolean | false | false | 인덱스가 없으면 관리할지 여부예요. |
| maxRetries | Integer | false | 1 | Elasticsearch 요청의 최대 재시도 횟수예요. 비활성화하려면 -1을 사용하세요. |
| retryBackoffInMs | Integer | false | 100 | Elasticsearch 요청 재시도 시 기다리는 기본 시간(밀리초)이에요. |
| maxRetryTimeInSec | Integer | false | 86400 | Elasticsearch 요청 재시도의 최대 재시도 시간 간격(초)이에요. |
| bulkEnabled | Boolean | false | false | 요청 수나 크기, 또는 일정 시간이 지난 뒤 쓰기 요청을 flush하도록 Elasticsearch 벌크 프로세서를 활성화할지 여부예요. |
| bulkActions | Integer | false | 1000 | Elasticsearch 벌크 요청 하나당 최대 액션 수예요. 비활성화하려면 -1을 사용하세요. |
| bulkSizeInMb | Integer | false | 5 | Elasticsearch 벌크 요청의 최대 크기(메가바이트)예요. 비활성화하려면 -1을 사용하세요. |
| bulkConcurrentRequests | Integer | false | 0 | 진행 중(in flight)인 Elasticsearch 벌크 요청의 최대 수예요. 기본값 0은 단일 요청 실행을 허용해요. 값이 1이면 새로운 벌크 요청을 누적하는 동안 동시 요청 1개를 실행할 수 있다는 뜻이에요. |
| bulkFlushIntervalInMs | Long | false | 1000 | 벌크 쓰기가 활성화됐을 때 보류 중인 쓰기를 flush하기 위해 기다리는 최대 시간이에요. -1 또는 0이면 예약된 flush가 비활성화돼요. |
| compressionEnabled | Boolean | false | false | Elasticsearch 요청 압축을 활성화할지 여부예요. |
| connectTimeoutInMs | Integer | false | 5000 | Elasticsearch 클라이언트 연결 타임아웃(밀리초)이에요. |
| connectionRequestTimeoutInMs | Integer | false | 1000 | Elasticsearch 연결 풀에서 연결을 가져오는 시간(밀리초)이에요. |
| connectionIdleTimeoutInMs | Integer | false | 5 | 읽기 타임아웃을 방지하는 유휴(idle) 연결 타임아웃이에요. |
| keyIgnore | Boolean | false | true | Elasticsearch 도큐먼트 _id를 만들 때 레코드 키를 무시할지 여부예요. primaryFields가 정의되면 커넥터는 페이로드에서 기본 필드를 추출해 도큐먼트 _id를 만들어요. primaryFields가 제공되지 않으면 elasticsearch가 무작위 도큐먼트 _id를 자동 생성해요. |
| primaryFields | String | false | "id" | 레코드 값에서 Elasticsearch 도큐먼트 _id를 만드는 데 사용하는 필드 이름의 쉼표 구분 순서 목록이에요. 이 목록이 단일 항목이면 필드는 문자열로 변환돼요. 필드가 2개 이상이면 생성된 _id는 필드 값들의 JSON 배열의 문자열 표현이에요. |
| nullValueAction | enum(IGNORE, DELETE, FAIL) | false | IGNORE | NULL 값을 가진 레코드를 어떻게 처리할지예요. 가능한 옵션은 IGNORE, DELETE, FAIL이에요. 기본은 메시지를 무시(IGNORE)해요. |
| malformedDocAction | enum(IGNORE, WARN, FAIL) | false | FAIL | 일부 변형(malformation)으로 Elasticsearch가 거부한 도큐먼트를 어떻게 처리할지예요. 가능한 옵션은 IGNORE, DELETE, FAIL이에요. 기본은 Elasticsearch 도큐먼트를 실패(FAIL) 처리해요. |
| stripNulls | Boolean | false | true | stripNulls가 false이면 elasticsearch _source에 빈 필드를 위해 'null'이 포함되고(예: {"foo": null}), 그렇지 않으면 null 필드는 제거돼요. |
| socketTimeoutInMs | Integer | false | 60000 | elasticsearch 응답을 읽기 위해 기다리는 소켓 타임아웃(밀리초)이에요. |
| typeName | String | false | "_doc" | 커넥터가 메시지를 쓰는 타입 이름이에요. Elasticsearch 6.2 이전 버전에서는 "_doc"가 아닌 유효한 타입 이름으로 명시적으로 설정해야 하고, 그 외에는 기본값을 그대로 두면 돼요. |
| indexNumberOfShards | int | false | 1 | 인덱스의 샤드 수예요. |
| indexNumberOfReplicas | int | false | 1 | 인덱스의 복제본 수예요. |
| username | String | false | " " (빈 문자열) | 커넥터가 Elasticsearch 클러스터에 연결하는 데 사용하는 사용자 이름이에요. username을 설정하면 password도 제공해야 해요. |
| password | String | false | " " (빈 문자열) | 커넥터가 Elasticsearch 클러스터에 연결하는 데 사용하는 비밀번호예요. username을 설정하면 password도 제공해야 해요. |
| ssl | ElasticSearchSslConfig | false | TLS 암호화 통신을 위한 구성이에요. | |
| compatibilityMode | enum(AUTO, ELASTICSEARCH, ELASTICSEARCH_7, OPENSEARCH) | false | AUTO | ElasticSearch 클러스터와의 호환 모드를 지정해요. AUTO 값은 사용할 올바른 호환 모드를 자동 감지하려 시도해요. 대상 클러스터가 ElasticSearch 7 이하를 실행 중이면 ELASTICSEARCH_7을 사용하세요. ElasticSearch 8 이상이면 ELASTICSEARCH를, OpenSearch이면 OPENSEARCH를 사용하세요. |
| token | String | false | " " (빈 문자열) | 커넥터가 ElasticSearch 클러스터에 연결하는 데 사용하는 토큰이에요. basic/token/apiKey 인증 모드 중 하나만 구성해야 해요. |
| apiKey | String | false | " " (빈 문자열) | 커넥터가 ElasticSearch 클러스터에 연결하는 데 사용하는 apiKey예요. basic/token/apiKey 인증 모드 중 하나만 구성해야 해요. |
| canonicalKeyFields | Boolean | false | false | JSON과 Avro의 키 필드를 정렬할지 여부예요. true로 설정하고 레코드 키 스키마가 JSON 또는 AVRO라면, 직렬화된 객체는 프로퍼티의 순서를 고려하지 않아요. |
| stripNonPrintableCharacters | Boolean | false | true | 도큐먼트에서 모든 비인쇄(non-printable) 문자를 제거할지 여부예요. true로 설정하면 모든 비인쇄 문자가 도큐먼트에서 제거돼요. |
| idHashingAlgorithm | enum(NONE, SHA256, SHA512) | false | NONE | 도큐먼트 ID에 사용할 해싱 알고리즘이에요. ElasticSearch _id의 512바이트 하드 리미트를 준수하는 데 유용해요. |
| conditionalIdHashing | Boolean | false | false | 이 옵션은 idHashingAlgorithm이 설정된 경우에만 동작해요. 활성화하면 ID가 512바이트보다 클 때만 해싱이 수행되고, 그렇지 않으면 어떤 경우든 각 도큐먼트에 해싱이 수행돼요. |
| copyKeyFields | Boolean | false | false | 메시지 키 스키마가 AVRO 또는 JSON이면 메시지 키 필드가 ElasticSearch 도큐먼트에 복사돼요. |
ElasticSearchSslConfig 구조 정의 (Definition of ElasticSearchSslConfig structure)
| 이름 | 타입 | 필수 | 기본값 | 설명 |
|---|---|---|---|---|
| enabled | Boolean | false | false | SSL/TLS를 활성화할지 여부예요. |
| hostnameVerification | Boolean | false | true | SSL 사용 시 노드 호스트네임을 검증할지 여부예요. |
| disableCertificateValidation | Boolean | false | true | 노드 인증서 검증을 비활성화할지 여부예요. 이 값을 바꾸는 건 매우 안전하지 않아서 프로덕션 환경에서는 사용하면 안 돼요. |
| truststorePath | String | false | " " (빈 문자열) | truststore 파일 경로예요. |
| truststorePassword | String | false | " " (빈 문자열) | truststore 비밀번호예요. |
| keystorePath | String | false | " " (빈 문자열) | keystore 파일 경로예요. |
| keystorePassword | String | false | " " (빈 문자열) | keystore 비밀번호예요. |
| cipherSuites | String | false | " " (빈 문자열) | SSL/TLS cipher suites예요. |
| protocols | String | false | "TLSv1.2" | 활성화된 SSL/TLS 프로토콜의 쉼표 구분 목록이에요. |
예제 (Example)
Elasticsearch sink 커넥터를 사용하기 전에 다음 방법 중 하나로 구성 파일을 만들어야 해요.
구성 (Configuration)
Elasticsearch 6.2 이후용
- JSON
{
"configs": {
"elasticSearchUrl": "http://localhost:9200",
"indexName": "my_index",
"username": "scooby",
"password": "doobie"
}
}
- YAML
configs:
elasticSearchUrl: "http://localhost:9200"
indexName: "my_index"
username: "scooby"
password: "doobie"
Elasticsearch 6.2 이전용
- JSON
{
"elasticSearchUrl": "http://localhost:9200",
"indexName": "my_index",
"typeName": "doc",
"username": "scooby",
"password": "doobie"
}
- YAML
configs:
elasticSearchUrl: "http://localhost:9200"
indexName: "my_index"
typeName: "doc"
username: "scooby"
password: "doobie"
사용법 (Usage)
- 단일 노드 Elasticsearch 클러스터를 시작해요.
docker run -p 9200:9200 -p 9300:9300 \
-e "discovery.type=single-node" \
docker.elastic.co/elasticsearch/elasticsearch:7.13.3
- Pulsar 서비스를 로컬 standalone 모드로 시작해요.
bin/pulsar standalone
NAR 파일이 connectors/pulsar-io-elastic-search-5.0.0-M2.nar 경로에 있는지 확인하세요.
-
다음 방법 중 하나로 Pulsar Elasticsearch 커넥터를 로컬 실행(local run) 모드로 시작해요.
-
앞에서 보여준 JSON 구성을 사용해요.
bin/pulsar-admin sinks localrun \
--archive $PWD/connectors/pulsar-io-elastic-search-5.0.0-M2.nar \
--tenant public \
--namespace default \
--name elasticsearch-test-sink \
--sink-config '{"elasticSearchUrl":"http://localhost:9200","indexName": "my_index","username": "scooby","password": "doobie"}' \
--inputs elasticsearch_test
- 앞에서 보여준 YAML 구성 파일을 사용해요.
bin/pulsar-admin sinks localrun \
--archive $PWD/connectors/pulsar-io-elastic-search-5.0.0-M2.nar \
--tenant public \
--namespace default \
--name elasticsearch-test-sink \
--sink-config-file $PWD/elasticsearch-sink.yml \
--inputs elasticsearch_test
- 토픽에 레코드를 발행해요.
bin/pulsar-client produce elasticsearch_test --messages "{\"a\":1}"
-
Elasticsearch에서 도큐먼트를 확인해요.
-
인덱스를 refresh해요.
curl -s http://localhost:9200/my_index/_refresh
- 도큐먼트를 검색해요.
curl -s http://localhost:9200/my_index/_search
앞서 발행한 레코드가 Elasticsearch에 정상적으로 쓰였는지 확인할 수 있어요.
{"took":2,"timed_out":false,"_shards":{"total":1,"successful":1,"skipped":0,"failed":0},"hits":{"total":{"value":1,"relation":"eq"},"max_score":1.0,"hits":[{"_index":"my_index","_type":"_doc","_id":"FSxemm8BLjG_iC0EeTYJ","_score":1.0,"_source":{"a":1}}]}}