Opensearch SQL 커넥터

Opensearch SQL 커넥터

Opensearch 커넥터는 Opensearch 엔진의 인덱스에 쓰는 것을 지원합니다. 이 문서는 Opensearch에 대해 SQL 쿼리를 실행할 수 있도록 Opensearch 커넥터를 설정하는 방법을 설명합니다.

이 커넥터는 DDL에 정의된 기본 키를 사용해 외부 시스템과 UPDATE/DELETE 메시지를 교환하는 upsert 모드로 동작할 수 있습니다.

DDL에 기본 키가 정의되지 않으면 커넥터는 외부 시스템과 INSERT만 메시지를 교환하는 append 모드로만 동작할 수 있습니다.

출처: 문서

본문

의존성

Flink 2.3 버전에 사용 가능한 커넥터는 아직 없습니다.

Opensearch 커넥터는 바이너리 배포판에 포함되지 않습니다. 클러스터 실행을 위해 링크하는 방법은 여기를 참조하세요.

Opensearch 테이블 생성 방법

다음 예시는 Opensearch sink 테이블을 생성하는 방법을 보여줍니다.

CREATE TABLE myUserTable (
  user_id STRING,
  user_name STRING,
  uv BIGINT,
  pv BIGINT,
  PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
  'connector' = 'opensearch',
  'hosts' = 'http://localhost:9200',
  'index' = 'users'
);

커넥터 옵션

Option Required Forwarded Default Type Description
connector required no (none) String 사용할 커넥터를 지정합니다. 유효한 값은: opensearch
hosts required yes (none) String 연결할 하나 이상의 Opensearch 호스트입니다. 예: 'http://host_name:9092;http://host_name:9093'.
index required yes (none) String 모든 레코드에 대한 Opensearch 인덱스입니다. 정적 인덱스(예: 'myIndex') 또는 동적 인덱스(예: 'index-{log_ts|yyyy-MM-dd}')가 될 수 있습니다. 자세한 내용은 아래 Dynamic Index 섹션을 참조하세요.
allow-insecure optional yes (none) Boolean HTTPS 엔드포인트에 대한 안전하지 않은 연결을 허용합니다(인증서 검증 비활성화).
document-id.key-delimiter optional yes _ String 복합 키의 구분자("_" 기본값)입니다. 예: "$"는 "KEY1$KEY2$KEY3"과 같은 ID를 만듭니다.
username optional yes (none) String Opensearch 인스턴스에 연결하는 데 사용되는 사용자 이름입니다. Opensearch에는 사전 번들된 보안 기능이 제공되며, guidelines에 따라 Opensearch 클러스터의 보안을 구성하여 비활성화할 수 있습니다.
password optional yes (none) String Opensearch 인스턴스에 연결하는 데 사용되는 비밀번호입니다. username이 구성되면 이 옵션도 비어 있지 않은 문자열로 구성해야 합니다.
sink.delivery-guarantee optional no AT_LEAST_ONCE String 커밋 시 선택적 전달 보장입니다. 유효한 값은: - EXACTLY_ONCE: 레코드는 장애 조치 시나리오에서도 exactly-once로만 전달됩니다. - AT_LEAST_ONCE: 레코드가 전달됨을 보장하지만 같은 레코드가 여러 번 전달될 수 있습니다. - NONE: 베스트 에포트 방식으로 레코드를 전달합니다.
sink.flush-on-checkpoint optional true Boolean checkpoint 시 플러시할지 여부입니다. 비활성화하면 sink는 checkpoint에서 Opensearch가 모든 보류 중인 작업 요청을 승인할 때까지 기다리지 않습니다. 따라서 sink는 작업 요청의 at-least-once 전달에 대한 강한 보장을 제공하지 않습니다.
sink.bulk-flush.max-actions optional yes 1000 Integer 벌크 요청당 버퍼링된 최대 액션 수입니다. 비활성화하려면 '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-size''sink.bulk-flush.max-actions''0'으로 설정하고 플러시 간격을 설정하면 버퍼링된 액션의 완전한 비동기 처리가 가능합니다.
sink.bulk-flush.backoff.strategy optional yes DISABLED String 일시적인 요청 오류로 인해 플러시 액션이 실패했을 때 재시도를 수행하는 방법을 지정합니다. 유효한 전략은: - DISABLED: 재시도하지 않습니다. 즉 첫 번째 요청 오류 후 실패합니다. - CONSTANT: 재시도 사이에 백오프 지연을 기다립니다. - EXPONENTIAL: 처음에 백오프 지연을 기다렸다가 재시도 사이에 지수적으로 증가합니다.
sink.bulk-flush.backoff.max-retries optional yes (none) Integer 최대 백오프 재시도 횟수입니다.
sink.bulk-flush.backoff.delay optional yes (none) Duration 각 백오프 시도 사이의 지연입니다. CONSTANT 백오프의 경우 단순히 재시도 사이의 지연입니다. EXPONENTIAL 백오프의 경우 초기 기본 지연입니다.
connection.path-prefix optional yes (none) String 모든 REST 통신에 추가할 접두사 문자열입니다. 예: '/v1'.
connection.request-timeout optional yes (none) Duration 연결 관리자에서 연결을 요청하는 타임아웃입니다.
connection.timeout optional yes (none) Duration 연결을 설정하는 타임아웃입니다.
socket.timeout optional yes (none) Duration 데이터를 기다리는 소켓 타임아웃(SO_TIMEOUT)입니다. 다시 말해 두 연속 데이터 패킷 사이의 최대 비활동 기간입니다.
format optional no json String Opensearch 커넥터는 포맷 지정을 지원합니다. 포맷은 유효한 json 문서를 생성해야 합니다. 기본적으로 내장 'json' 포맷을 사용합니다. 자세한 내용은 JSON Format 페이지를 참조하세요.

기능

키 처리

Opensearch sink는 기본 키가 정의되었는지에 따라 upsert 모드 또는 append 모드로 동작할 수 있습니다. 기본 키가 정의되면 Opensearch sink는 UPDATE/DELETE 메시지를 포함하는 쿼리를 소비할 수 있는 upsert 모드로 동작합니다. 기본 키가 정의되지 않으면 Opensearch sink는 INSERT만 메시지를 포함하는 쿼리만 소비할 수 있는 append 모드로 동작합니다.

Opensearch 커넥터에서 기본 키는 Opensearch 문서 ID를 계산하는 데 사용되며, 이는 최대 512바이트의 문자열입니다. 공백을 포함할 수 없습니다. Opensearch 커넥터는 document-id.key-delimiter로 지정된 키 구분자를 사용해 DDL에 정의된 순서대로 모든 기본 키 필드를 연결하여 모든 행에 대한 문서 ID 문자열을 생성합니다. BYTES, ROW, ARRAY, MAP 등 일부 타입은 좋은 문자열 표현이 없으므로 기본 키 필드로 허용되지 않습니다. 기본 키가 지정되지 않으면 Opensearch가 문서 ID를 자동으로 생성합니다.

PRIMARY KEY 구문에 대한 자세한 내용은 CREATE TABLE DDL을 참조하세요.

동적 인덱스

Opensearch sink는 정적 인덱스와 동적 인덱스를 모두 지원합니다.

정적 인덱스를 사용하려면 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 스트림만 지원할 수 있습니다.

데이터 타입 매핑

Opensearch는 문서를 JSON 문자열로 저장합니다. 따라서 데이터 타입 매핑은 Flink 데이터 타입과 JSON 데이터 타입 사이입니다. Flink는 Opensearch 커넥터에 내장 'json' 포맷을 사용합니다. 더 자세한 타입 매핑은 JSON Format 페이지를 참조하세요.

더 알아보기 (Learn more)