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}}]}}

더 알아보기 (Learn more)