Elasticsearch SQL 커넥터
Elasticsearch SQL 커넥터 (Elasticsearch SQL Connector)
Sink: Batch, Streaming Append & Upsert 모드
Elasticsearch 커넥터는 Elasticsearch 엔진의 인덱스에 데이터를 쓸 수 있게 해줍니다. 이 문서는 Elasticsearch 커넥터를 설정하고 Elasticsearch를 대상으로 SQL 쿼리를 실행하는 방법을 설명합니다.
출처: 문서
본문
커넥터는 DDL에 정의된 기본 키(primary key)를 사용해 외부 시스템과 UPDATE/DELETE 메시지를 교환하는 upsert 모드로 동작할 수 있습니다.
DDL에 기본 키가 정의되지 않으면 커넥터는 외부 시스템과 INSERT 메시지만 교환하는 append 모드로만 동작할 수 있습니다.
의존성 (Dependencies)
아직 Flink 2.3용 커넥터를 사용할 수 없습니다.
Elasticsearch 커넥터는 바이너리 배포판에 포함되어 있지 않습니다. 클러스터 실행을 위해 연결하는 방법은 여기를 참고하세요.
Elasticsearch 테이블을 만드는 방법
다음 예제는 Elasticsearch 싱크 테이블을 만드는 방법을 보여줍니다.
CREATE TABLE myUserTable (
user_id STRING,
user_name STRING,
uv BIGINT,
pv BIGINT,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'elasticsearch-7',
'hosts' = 'http://localhost:9200',
'index' = 'users'
);
커넥터 옵션 (Connector Options)
| 옵션 | 필수 | 전달됨 | 기본값 | 유형 | 설명 |
|---|---|---|---|---|---|
| connector | required | no | (없음) | String | 사용할 커넥터를 지정합니다. 유효한 값은 elasticsearch-6(6.x 클러스터), elasticsearch-7(7.x 클러스터), elasticsearch-8(8.x 클러스터)입니다. |
| hosts | required | yes | (없음) | String | 연결할 Elasticsearch 호스트 목록입니다. 예: 'http://host_name:9092;http://host_name:9093'. |
| index | required | yes | (없음) | String | 각 레코드에 대한 Elasticsearch 인덱스입니다. 정적 인덱스(예: 'myIndex') 또는 동적 인덱스(예: 'index-{log_ts|yyyy-MM-dd}')가 될 수 있습니다. |
| document-type | 6.x에서 required | 6.x에서 yes | (없음) | String | Elasticsearch 문서 타입입니다. elasticsearch-7과 elasticsearch-8에서는 더 이상 필요하지 않습니다. |
| document-id.key-delimiter | optional | yes | _ | String | 복합 키의 구분자(기본 "_")입니다. 예를 들어 "$"이면 ID가 "KEY1$KEY2$KEY3"가 됩니다. |
| username | optional | yes | (없음) | String | Elasticsearch 인스턴스에 연결하는 데 사용할 사용자 이름입니다. |
| password | optional | yes | (없음) | String | Elasticsearch 인스턴스에 연결하는 데 사용할 비밀번호입니다. username이 설정되면 이 옵션도 비어 있지 않은 문자열로 설정해야 합니다. |
| failure-handler | optional | yes | fail | String | Elasticsearch 요청이 실패할 때의 실패 처리 전략입니다. fail(예외 발생·작업 실패), ignore(무시하고 요청 폐기), retry-rejected(큐 용량 포화로 실패한 요청 재추가), 사용자 정의 클래스 이름(ActionRequestFailureHandler 하위 클래스)을 지원합니다. |
| sink.flush-on-checkpoint | optional | true | Boolean | 체크포인트에서 플러시할지 여부입니다. 비활성화하면 증거가 약해집니다. | |
| sink.bulk-flush.max-actions | optional | yes | 1000 | Integer | 벌크 요청당 버퍼링할 최대 액션 수입니다. '0'으로 비활성화할 수 있습니다(elasticsearch-8에서는 0보다 커야 함). |
| sink.bulk-flush.max-size | optional | yes | 2mb | MemorySize | 벌크 요청당 버퍼링 액션의 최대 메모리 크기입니다. MB 단위여야 합니다. '0'으로 비활성화 가능합니다. |
| sink.bulk-flush.interval | optional | yes | 1s | Duration | 버퍼링된 액션을 플러시할 간격입니다. '0'으로 비활성화 가능합니다. |
| sink.bulk-flush.max-buffered-actions | optional | yes | 10000 | Integer | 싱크에 버퍼링될 수 있는 최대 레코드 수이며 sink.bulk-flush.max-actions보다 커야 합니다. elasticsearch-8에서만 지원됩니다. |
| sink.bulk-flush.max-in-flight-actions | optional | yes | 50 | Integer | 허용되는 최대 in-flight 요청 수입니다. elasticsearch-8에서만 지원됩니다. |
| sink.bulk-flush.backoff.strategy | optional | yes | DISABLED | String | 임시 요청 오류로 플러시가 실패할 때 재시도하는 방법입니다. DISABLED, CONSTANT, EXPONENTIAL을 지원합니다. elasticsearch-8에서는 지원되지 않습니다. |
| sink.bulk-flush.backoff.max-retries | optional | yes | (없음) | Integer | 백오프 재시도의 최대 횟수입니다. elasticsearch-8에서는 지원되지 않습니다. |
| sink.bulk-flush.backoff.delay | optional | yes | (없음) | Duration | 각 백오프 시도 사이의 지연입니다. elasticsearch-8에서는 지원되지 않습니다. |
| ssl.certificate-fingerprint | optional | yes | (없음) | String | HTTPS 연결 검증에 사용되는 HTTP CA 인증서 SHA-256 핑거프린트입니다. elasticsearch-8에서만 지원됩니다. |
| connection.path-prefix | optional | yes | (없음) | String | 모든 REST 통신에 추가되는 접두사 문자열입니다. 예: '/v1'. |
| format | optional | no | json | String | Elasticsearch 커넥터는 포맷 지정을 지원합니다. 포맷은 유효한 JSON 문서를 생성해야 하며 기본적으로 내장 'json' 포맷을 사용합니다. |
기능 (Features)
키 처리 (Key Handling)
Elasticsearch 싱크는 기본 키 정의 여부에 따라 upsert 모드 또는 append 모드로 동작합니다. 기본 키가 정의되면 upsert 모드로 동작해 UPDATE/DELETE 메시지를 포함한 쿼리를 사용할 수 있습니다. 기본 키가 정의되지 않으면 INSERT 메시지만 포함된 쿼리를 사용할 수 있는 append 모드로 동작합니다.
Elasticsearch 커넥터에서 기본 키는 Elasticsearch 문서 ID를 계산하는 데 사용되며, 이는 최대 512바이트의 문자열로 공백을 가질 수 없습니다. 커넥터는 DDL에 정의된 순서대로 모든 기본 키 필드를 document-id.key-delimiter가 지정하는 구분자로 연결해 각 행에 대한 문서 ID 문자열을 생성합니다. BYTES, ROW, ARRAY, MAP 등 좋은 문자열 표현이 없는 일부 타입은 기본 키 필드로 허용되지 않습니다. 기본 키가 지정되지 않으면 Elasticsearch가 문서 ID를 자동 생성합니다.
PRIMARY KEY 문법에 대한 자세한 내용은 CREATE TABLE DDL을 참고하세요.
동적 인덱스 (Dynamic Index)
Elasticsearch 싱크는 정적 인덱스와 동적 인덱스를 모두 지원합니다.
정적 인덱스를 원한다면 index 옵션 값을 일반 문자열(예: 'myusers')로 지정하면 모든 레코드가 "myusers" 인덱스에 일관되게 기록됩니다.
동적 인덱스를 원한다면 {field_name}을 사용해 레코드의 필드 값을 참조해 대상 인덱스를 동적으로 생성할 수 있습니다. '{field_name|date_format_string}'을 사용하면 TIMESTAMP/DATE/TIME 타입 필드 값을 date_format_string이 지정하는 포맷으로 변환할 수 있습니다. date_format_string은 Java의 DateTimeFormatter와 호환됩니다. 예를 들어 옵션 값이 'myusers-{log_ts|yyyy-MM-dd}'이면, log_ts 필드 값이 2020-03-27 12:25:55인 레코드는 "myusers-2020-03-27" 인덱스에 기록됩니다.
또한 '{now()|date_format_string}'을 사용해 현재 시스템 시간을 date_format_string이 지정하는 포맷으로 변환할 수 있습니다. now()의 해당 시간 타입은 TIMESTAMP_WITH_LTZ입니다. 시스템 시간을 문자열로 포맷할 때는 세션에서 table.local-time-zone으로 설정한 시간대가 사용됩니다. NOW(), now(), CURRENT_TIMESTAMP, current_timestamp를 사용할 수 있습니다.
참고: 현재 시스템 시간으로 생성한 동적 인덱스를 사용할 때, changelog 스트림에서는 같은 기본 키를 가진 레코드가 같은 인덱스 이름을 생성한다는 보장이 없습니다. 따라서 시스템 시간 기반의 동적 인덱스는 append-only 스트림만 지원할 수 있습니다.
데이터 타입 매핑 (Data Type Mapping)
Elasticsearch는 문서를 JSON 문자열로 저장하므로, 데이터 타입 매핑은 Flink 데이터 타입과 JSON 데이터 타입 사이의 매핑입니다. Flink는 Elasticsearch 커넥터에 내장 'json' 포맷을 사용합니다. 자세한 타입 매핑은 JSON Format 페이지를 참고하세요.