HBase SQL 커넥터
HBase SQL 커넥터
HBase 커넥터는 HBase 클러스터에서 데이터를 읽고 쓸 수 있게 해줘요. 이 문서는 HBase 커넥터를 설정해 HBase에 대해 SQL 쿼리를 실행하는 방법을 설명해요.
출처: 문서
본문
HBase는 항상 upsert 모드로 동작하며, DDL에 정의된 기본 키를 사용해 외부 시스템과 changelog 메시지를 교환해요. 기본 키는 HBase rowkey 필드에 정의되어야 해요(rowkey 필드는 선언되어야 함). PRIMARY KEY 절이 선언되지 않으면, HBase 커넥터는 기본적으로 rowkey를 기본 키로 취급해요.
의존성 (Dependencies)
현재 Flink 2.3 버전용 커넥터는 아직 제공되지 않아요.
HBase 커넥터는 바이너리 배포판의 일부가 아니에요. 클러스터 실행을 위한 링크 방법은 여기를 참조하세요.
HBase 테이블 사용 방법 (How to use HBase table)
HBase 테이블의 모든 컬럼 패밀리(column family)는 ROW 타입으로 선언되어야 해요. 필드 이름은 컬럼 패밀리 이름에 매핑되고, 중첩 필드 이름은 컬럼 퀄리파이어(column qualifier) 이름에 매핑돼요. 스키마에 모든 패밀리와 퀄리파이어를 선언할 필요는 없으며, 사용자는 쿼리에서 사용되는 것만 선언할 수 있어요. ROW 타입 필드를 제외한 단일 원자(atomic) 타입 필드(예: STRING, BIGINT)는 HBase rowkey로 인식돼요. rowkey 필드는 임의의 이름이 될 수 있지만, 예약 키워드라면 백틱을 사용해 인용해야 해요.
-- register the HBase table 'mytable' in Flink SQL
CREATE TABLE hTable (
rowkey INT,
family1 ROW<q1 INT>,
family2 ROW<q2 STRING, q3 BIGINT>,
family3 ROW<q4 DOUBLE, q5 BOOLEAN, q6 STRING>,
PRIMARY KEY (rowkey) NOT ENFORCED
) WITH (
'connector' = 'hbase-2.2',
'table-name' = 'mytable',
'zookeeper.quorum' = 'localhost:2181'
);
-- use ROW(...) construction function construct column families and write data into the HBase table.
-- assuming the schema of "T" is [rowkey, f1q1, f2q2, f2q3, f3q4, f3q5, f3q6]
INSERT INTO hTable
SELECT rowkey, ROW(f1q1), ROW(f2q2, f2q3), ROW(f3q4, f3q5, f3q6) FROM T;
-- scan data from the HBase table
SELECT rowkey, family1, family3.q4, family3.q6 FROM hTable;
-- temporal join the HBase table as a dimension table
SELECT * FROM myTopic
LEFT JOIN hTable FOR SYSTEM_TIME AS OF myTopic.proctime
ON myTopic.key = hTable.rowkey;
사용 가능한 메타데이터 (Available Metadata)
다음 커넥터 메타데이터는 테이블 정의에서 메타데이터 컬럼으로 접근할 수 있어요.
R/W 컬럼은 메타데이터 필드가 읽기 가능(R) 및/또는 쓰기 가능(W)인지 정의해요.
읽기 전용 컬럼은 INSERT INTO 연산 중에 제외되도록 VIRTUAL로 선언해야 해요.
| 키(Key) | 데이터 타입(Data Type) | 설명(Description) | R/W |
|---|---|---|---|
timestamp |
TIMESTAMP_LTZ(3) NOT NULL |
HBase 변형(mutation)에 대한 타임스탬프. | W |
ttl |
BIGINT NOT NULL |
HBase 변형에 대한 TTL(Time-to-live), 밀리초 단위. | W |
커넥터 옵션 (Connector Options)
| 옵션(Option) | 필수 | 전달(Forwarded) | 기본값(Default) | 타입(Type) | 설명(Description) |
|---|---|---|---|---|---|
| connector | required | no | (none) | String | 어떤 커넥터를 사용할지 지정해요. 유효한 값: - hbase-2.2: HBase 2.2.x 클러스터에 연결 |
| table-name | required | yes | (none) | String | 연결할 HBase 테이블의 이름. 기본적으로 테이블은 'default' 네임스페이스에 있어요. 테이블에 지정된 네임스페이스를 할당하려면 'namespace:table'을 사용해야 해요. |
| zookeeper.quorum | required | yes | (none) | String | HBase Zookeeper 쿼럼. |
| zookeeper.znode.parent | optional | yes | /hbase | String | HBase 클러스터의 Zookeeper 루트 디렉터리. |
| null-string-literal | optional | yes | null | String | 문자열 필드의 null 값 표현. HBase 소스와 싱크는 문자열 타입을 제외한 모든 타입에 대해 빈 바이트를 null 값으로 인코딩/디코딩해요. |
| sink.buffer-flush.max-size | optional | yes | 2mb | MemorySize | 쓰기 옵션. 각 쓰기 요청에 대해 버퍼링되는 행의 최대 메모리 크기. HBase 데이터베이스에 데이터를 쓰는 성능을 개선할 수 있지만 지연 시간을 늘릴 수도 있어요. '0'으로 설정해 비활성화할 수 있어요. |
| sink.buffer-flush.max-rows | optional | yes | 1000 | Integer | 쓰기 옵션. 각 쓰기 요청에 대해 버퍼링할 행의 최대 수. 데이터를 쓰는 성능을 개선할 수 있지만 지연 시간을 늘릴 수도 있어요. '0'으로 설정해 비활성화할 수 있어요. |
| sink.buffer-flush.interval | optional | yes | 1s | Duration | 쓰기 옵션. 버퍼링된 행을 플러시할 간격. 데이터를 쓰는 성능을 개선할 수 있지만 지연 시간을 늘릴 수도 있어요. '0'으로 설정해 비활성화할 수 있어요. 참고: 'sink.buffer-flush.max-size'와 'sink.buffer-flush.max-rows'를 '0'으로 설정하고 플러시 간격을 설정하면 버퍼링된 액션을 완전히 비동기로 처리할 수 있어요. |
| sink.ignore-null-value | optional | yes | false | Boolean | 쓰기 옵션. null 값을 무시할지 여부. |
| sink.parallelism | optional | no | (none) | Integer | HBase 싱크 연산자의 병렬도를 정의해요. 기본적으로 병렬도는 업스트림 체인 연산자와 동일한 병렬도를 사용해 프레임워크에 의해 결정돼요. |
| lookup.async | optional | no | false | Boolean | 비동기 lookup이 활성화되었는지 여부. true면 lookup은 비동기가 돼요. 참고: async는 hbase-2.2 커넥터만 지원해요. |
| lookup.cache | optional | yes | NONE | Enum 가능한 값: NONE, PARTIAL | lookup 테이블의 캐시 전략. 현재 NONE(캐시 없음)과 PARTIAL(외부 데이터베이스의 lookup 연산에서 항목 캐싱)을 지원해요. |
| lookup.partial-cache.max-rows | optional | yes | (none) | Long | lookup 캐시의 최대 행 수. 이 값을 초과하면 가장 오래된 행이 만료돼요. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"이어야 해요. |
| lookup.partial-cache.expire-after-write | optional | yes | (none) | Duration | 캐시에 쓴 후 lookup 캐시의 각 행의 최대 TTL. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"이어야 해요. |
| lookup.partial-cache.expire-after-access | optional | yes | (none) | Duration | 캐시의 항목에 접근한 후 lookup 캐시의 각 행의 최대 TTL. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"이어야 해요. |
| lookup.partial-cache.caching-missing-key | optional | yes | true | Boolean | lookup 키가 테이블의 어떤 행과도 일치하지 않을 때 캐시에 빈 값을 저장할지 여부. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"이어야 해요. |
| lookup.max-retries | optional | yes | 3 | Integer | lookup 데이터베이스가 실패할 경우 최대 재시도 횟수. |
| properties.* | optional | no | (none) | String | 임의의 HBase 구성을 설정하고 전달할 수 있어요. 접미사 이름은 HBase Configuration 문서에 정의된 구성 키와 일치해야 해요. Flink는 "properties." 키 접두사를 제거하고 변환된 키와 값을 기본 HBaseClient에 전달해요. 예를 들어 kerberos 인증 파라미터 'properties.hbase.security.authentication' = 'kerberos'를 추가할 수 있어요. |
비권장 옵션 (Deprecated Options)
이 비권장 옵션들은 위에 나열된 새 옵션으로 대체되었으며 결국 제거될 거예요. 새 옵션을 먼저 사용하는 것을 고려하세요.
| 옵션(Option) | 필수 | 전달(Forwarded) | 기본값(Default) | 타입(Type) | 설명(Description) |
|---|---|---|---|---|---|
| lookup.cache.max-rows | optional | yes | (none) | Integer | "lookup.cache" = "PARTIAL"로 설정하고 "lookup.partial-cache.max-rows"를 사용하세요. |
| lookup.cache.ttl | optional | yes | (none) | Duration | "lookup.cache" = "PARTIAL"로 설정하고 "lookup.partial-cache.expire-after-write"를 사용하세요. |
데이터 타입 매핑 (Data Type Mapping)
HBase는 모든 데이터를 바이트 배열(byte arrays)로 저장해요. 데이터는 읽기와 쓰기 연산 중에 직렬화/역직렬화되어야 해요.
직렬화와 역직렬화 시 Flink HBase 커넥터는 HBase(Hadoop)에서 제공하는 유틸리티 클래스 org.apache.hadoop.hbase.util.Bytes를 사용해 Flink 데이터 타입을 바이트 배열로 변환하고 그 반대로도 변환해요.
Flink HBase 커넥터는 문자열 타입을 제외한 모든 데이터 타입에 대해 null 값을 빈 바이트로 인코딩하고, 빈 바이트를 null 값으로 디코딩해요. 문자열 타입의 경우 null 리터럴은 null-string-literal 옵션에 의해 결정돼요.
데이터 타입 매핑은 다음과 같아요:
| Flink SQL 타입 | HBase 변환 |
|---|---|
CHAR / VARCHAR / STRING |
byte[] toBytes(String s) String toString(byte[] b) |
BOOLEAN |
byte[] toBytes(boolean b) boolean toBoolean(byte[] b) |
BINARY / VARBINARY |
byte[] 그대로 반환. |
DECIMAL |
byte[] toBytes(BigDecimal v) BigDecimal toBigDecimal(byte[] b) |
TINYINT |
new byte[] { val } bytes[0] // bytes의 첫 번째 유일한 바이트를 반환 |
SMALLINT |
byte[] toBytes(short val) short toShort(byte[] bytes) |
INT |
byte[] toBytes(int val) int toInt(byte[] bytes) |
BIGINT |
byte[] toBytes(long val) long toLong(byte[] bytes) |
FLOAT |
byte[] toBytes(float val) float toFloat(byte[] bytes) |
DOUBLE |
byte[] toBytes(double val) double toDouble(byte[] bytes) |
DATE |
epoch 이후의 일 수를 int 값으로 저장. |
TIME |
하루 중 밀리초 수를 int 값으로 저장. |
TIMESTAMP |
epoch 이후의 밀리초를 long 값으로 저장. |
ARRAY |
지원하지 않음 |
MAP / MULTISET |
지원하지 않음 |
ROW |
지원하지 않음 |