Apache Kafka SQL 커넥터
Apache Kafka SQL 커넥터 (Apache Kafka SQL Connector)
스캔 소스: 무한(Unbounded) / 싱크: 스트리밍 Append 모드
Kafka 커넥터는 Kafka 토픽에서 데이터를 읽고 Kafka 토픽으로 데이터를 쓸 수 있게 해줍니다.
출처: 문서
본문
의존성 (Dependencies)
현재 Flink 2.3 버전용 커넥터는 아직 없습니다.
Kafka 커넥터는 바이너리 배포의 일부가 아닙니다. 클러스터 실행을 위해 연결하는 방법은 여기를 참고하세요.
Kafka 테이블 생성 방법 (How to create a Kafka table)
아래 예시는 Kafka 테이블을 만드는 방법을 보여줍니다:
CREATE TABLE KafkaTable (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING,
`ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp'
) WITH (
'connector' = 'kafka',
'topic' = 'user_behavior',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'scan.startup.mode' = 'earliest-offset',
'format' = 'csv'
)
사용 가능한 메타데이터 (Available Metadata)
다음 커넥터 메타데이터는 테이블 정의에서 메타데이터 컬럼으로 접근할 수 있습니다.
R/W 컬럼은 메타데이터 필드가 읽을 수 있는지(R) 및/또는 쓸 수 있는지(W)를 정의합니다. 읽기 전용 컬럼은 INSERT INTO 연산 중 제외되도록 VIRTUAL로 선언해야 합니다.
| 키 | 데이터 타입 | 설명 | R/W |
|---|---|---|---|
topic |
STRING NOT NULL |
Kafka 레코드의 토픽 이름. | R/W |
partition |
INT NOT NULL |
Kafka 레코드의 파티션 ID. | R |
headers |
MAP NOT NULL |
원시 바이트 맵으로서의 Kafka 레코드의 헤더. | R/W |
leader-epoch |
INT NULL |
사용 가능한 경우 Kafka 레코드의 리더 epoch. | R |
offset |
BIGINT NOT NULL |
파티션 내 Kafka 레코드의 오프셋. | R |
timestamp |
TIMESTAMP_LTZ(3) NOT NULL |
Kafka 레코드의 타임스탬프. | R/W |
timestamp-type |
STRING NOT NULL |
Kafka 레코드의 타임스탬프 유형. "NoTimestampType", "CreateTime"(메타데이터 쓰기 시에도 설정됨) 또는 "LogAppendTime" 중 하나. | R |
확장된 CREATE TABLE 예시는 이러한 메타데이터 필드를 노출하는 문법을 보여줍니다:
CREATE TABLE KafkaTable (
`event_time` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp',
`partition` BIGINT METADATA VIRTUAL,
`offset` BIGINT METADATA VIRTUAL,
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING
) WITH (
'connector' = 'kafka',
'topic' = 'user_behavior',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'scan.startup.mode' = 'earliest-offset',
'format' = 'csv'
);
포맷 메타데이터 (Format Metadata)
커넥터는 읽기를 위한 값 포맷의 메타데이터를 노출할 수 있습니다. 포맷 메타데이터 키는 'value.' 접두사로 시작합니다.
다음 예시는 Kafka와 Debezium 메타데이터 필드 둘 다에 접근하는 방법을 보여줍니다:
CREATE TABLE KafkaTable (
`event_time` TIMESTAMP_LTZ(3) METADATA FROM 'value.source.timestamp' VIRTUAL, -- from Debezium format
`origin_table` STRING METADATA FROM 'value.source.table' VIRTUAL, -- from Debezium format
`partition_id` BIGINT METADATA FROM 'partition' VIRTUAL, -- from Kafka connector
`offset` BIGINT METADATA VIRTUAL, -- from Kafka connector
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING
) WITH (
'connector' = 'kafka',
'topic' = 'user_behavior',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'scan.startup.mode' = 'earliest-offset',
'value.format' = 'debezium-json'
);
커넥터 옵션 (Connector Options)
| 옵션 | 필수 | Forwarded | 기본값 | 타입 | 설명 |
|---|---|---|---|---|---|
| connector | required | no | (none) | String | 사용할 커넥터. Kafka에는 'kafka'를 사용합니다. |
| topic | optional | yes | (none) | String | 테이블을 소스로 사용할 때 읽을 토픽 이름, 또는 싱크로 사용할 때 쓸 토픽. 소스에는 'topic-1;topic-2'처럼 세미콜론으로 토픽을 구분해 토픽 목록을 지원합니다. "topic-pattern"과 "topic" 중 하나만 지정할 수 있습니다. 싱크의 경우 토픽 이름이 데이터를 쓸 토픽입니다. 싱크에도 토픽 목록을 지원합니다. 제공된 토픽 목록은 topic 메타데이터 컬럼의 유효 값 허용 목록으로 취급됩니다. 목록이 제공되면 싱크 테이블의 'topic' 메타데이터 컬럼은 쓰기 가능하며 반드시 지정되어야 합니다. |
| topic-pattern | optional | yes | (none) | String | 읽거나 쓸 토픽 이름 패턴의 정규식. 지정된 정규식과 일치하는 이름의 모든 토픽은 작업이 시작되면 컨슈머가 구독합니다. 싱크의 경우 topic 메타데이터 컬럼이 쓰기 가능하고, 제공되어야 하며, topic-pattern 정규식과 일치해야 합니다. "topic-pattern"과 "topic" 중 하나만 지정할 수 있습니다. |
| properties.bootstrap.servers | required | yes | (none) | String | 콤마로 구분된 Kafka 브로커 목록. |
| properties.group.id | source에는 optional, sink에는 해당 없음 | yes | (none) | String | Kafka 소스의 컨슈머 그룹 ID. 그룹 ID가 지정되지 않으면 자동 생성된 "KafkaSource-{tableIdentifier}" ID가 사용됩니다. |
| properties.* | optional | no | (none) | String | 임의의 Kafka 구성을 설정·전달할 수 있습니다. 접미사 이름은 Kafka 구성 문서에 정의된 구성 키와 일치해야 합니다. Flink는 "properties." 키 접두사를 제거하고 변환된 키와 값을 기본 KafkaClient에 전달합니다. 예를 들어 'properties.allow.auto.create.topics' = 'false'로 자동 토픽 생성을 비활성화할 수 있습니다. 그러나 'auto.offset.reset'처럼 Flink가 재정의하므로 설정할 수 없는 구성도 있습니다. |
| format | required | no | (none) | String | Kafka 메시지의 값 부분을 역직렬화·직렬화하는 데 사용되는 포맷. 자세한 내용과 추가 포맷 옵션은 formats 페이지를 참고하세요. 이 옵션 또는 'value.format' 옵션 중 하나가 필요합니다. |
| key.format | optional | no | (none) | String | Kafka 메시지의 키 부분을 역직렬화·직렬화하는 데 사용되는 포맷. 자세한 내용은 formats 페이지를 참고하세요. 키 포맷이 정의되면 'key.fields' 옵션도 필요합니다. 그렇지 않으면 Kafka 레코드의 키가 비게 됩니다. |
| key.fields | optional | no | [] | List<String> | 키 포맷의 데이터 타입을 구성하는 테이블 스키마의 물리적 컬럼 명시적 목록을 정의합니다. 기본적으로 이 목록은 비어 있어 키가 정의되지 않습니다. 목록은 'field1;field2'처럼 보여야 합니다. |
| key.fields-prefix | optional | no | (none) | String | 값 포맷의 필드와 이름 충돌을 피하기 위해 키 포맷의 모든 필드에 대한 커스텀 접두사를 정의합니다. 기본적으로 접두사는 비어 있습니다. 커스텀 접두사가 정의되면 테이블 스키마와 'key.fields' 모두 접두사가 붙은 이름으로 동작합니다. 키 포맷의 데이터 타입을 구성할 때 접두사가 제거되고 키 포맷 안에서 접두사가 없는 이름이 사용됩니다. 이 옵션을 사용하려면 'value.fields-include'가 'EXCEPT_KEY'로 설정되어야 합니다. |
| value.format | required | no | (none) | String | Kafka 메시지의 값 부분을 역직렬화·직렬화하는 데 사용되는 포맷. 자세한 내용은 formats 페이지를 참고하세요. 이 옵션 또는 'format' 옵션 중 하나가 필요합니다. |
| value.fields-include | optional | no | ALL | Enum ([ALL, EXCEPT_KEY]) | 값 포맷의 데이터 타입에서 키 컬럼을 처리하는 전략을 정의합니다. 기본적으로 테이블 스키마의 'ALL' 물리적 컬럼이 값 포맷에 포함되며, 이는 키 컬럼이 키 포맷과 값 포맷 둘 다의 데이터 타입에 나타남을 의미합니다. |
| scan.startup.mode | optional | yes | group-offsets | Enum | Kafka 컨슈머의 시작 모드. 유효한 값은 'earliest-offset', 'latest-offset', 'group-offsets', 'timestamp', 'specific-offsets'입니다. 자세한 내용은 아래 Start Reading Position을 참고하세요. |
| scan.startup.specific-offsets | optional | yes | (none) | String | 'specific-offsets' 시작 모드의 경우 각 파티션의 오프셋을 지정합니다. 예: 'partition:0,offset:42;partition:1,offset:300'. |
| scan.startup.timestamp-millis | optional | yes | (none) | Long | 'timestamp' 시작 모드의 경우 사용되는 지정된 epoch 타임스탬프(밀리초)부터 시작합니다. |
| scan.bounded.mode | optional | no | unbounded | Enum | Kafka 컨슈머의 유한(bounded) 모드. 유효한 값은 'latest-offset', 'group-offsets', 'timestamp', 'specific-offsets'입니다. 자세한 내용은 아래 Bounded Ending Position을 참고하세요. |
| scan.bounded.specific-offsets | optional | yes | (none) | String | 'specific-offsets' 유한 모드의 경우 각 파티션의 오프셋을 지정합니다. 예: 'partition:0,offset:42;partition:1,offset:300'. 파티션에 대한 오프셋이 제공되지 않으면 해당 파티션에서 소비하지 않습니다. |
| scan.bounded.timestamp-millis | optional | yes | (none) | Long | 'timestamp' 유한 모드의 경우 사용되는 지정된 epoch 타임스탬프(밀리초)에서 끝납니다. |
| scan.topic-partition-discovery.interval | optional | yes | 5 minutes | Duration | 컨슈머가 동적으로 생성된 Kafka 토픽과 파티션을 주기적으로 발견하는 간격. 이 기능을 비활성화하려면 'scan.topic-partition-discovery.interval' 값을 0으로 명시적으로 설정해야 합니다. |
| scan.parallelism | optional | no | (none) | Integer | Kafka 소스 연산자의 병렬도를 정의합니다. 설정하지 않으면 전역 기본 병렬도가 사용됩니다. |
| sink.partitioner | optional | yes | 'default' | String | Flink의 파티션에서 Kafka의 파티션으로의 출력 파티셔닝. 유효한 값은 |
default: 레코드를 파티셔닝하기 위해 kafka 기본 파티셔너 사용.fixed: 각 Flink 파티션은 최대 하나의 Kafka 파티션에 들어감.round-robin: Flink 파티션이 Kafka 파티션에 sticky round-robin으로 분산됨. 레코드 키가 지정되지 않은 경우에만 동작.- 커스텀
FlinkKafkaPartitioner하위 클래스: 예:'org.mycompany.MyPartitioner'. 자세한 내용은 아래 Sink Partitioning을 참고하세요. | | sink.semantic | optional | no | at-least-once | String | 더 이상 사용되지 않음(Deprecated):sink.delivery-guarantee를 사용하세요. | | sink.delivery-guarantee | optional | no | at-least-once | String | Kafka 싱크의 전달 의미론을 정의합니다. 유효한 값은'at-least-once','exactly-once','none'입니다. 자세한 내용은 Consistency guarantees를 참고하세요. | | sink.transactional-id-prefix | optional | yes | (none) | String | 전달 보장이'exactly-once'로 구성되면 이 값을 설정해야 하며, 열리는 모든 Kafka 트랜잭션의 식별자 접두사로 사용됩니다. | | sink.parallelism | optional | no | (none) | Integer | Kafka 싱크 연산자의 병렬도를 정의합니다. 기본적으로 병렬도는 업스트림 체인 연산자의 병렬도와 동일하게 프레임워크가 결정합니다. |
기능 (Features)
키와 값 포맷 (Key and Value Formats)
Kafka 레코드의 키와 값 부분 모두 주어진 포맷 중 하나를 사용해 원시 바이트로 직렬화·역직렬화할 수 있습니다.
값 포맷 (Value Format)
Kafka 레코드에서 키는 선택 사항이므로, 다음 문은 키 포맷 없이 구성된 값 포맷으로 레코드를 읽고 씁니다. 'format' 옵션은 'value.format'의 동의어입니다. 모든 포맷 옵션은 포맷 식별자로 접두사가 붙습니다.
CREATE TABLE KafkaTable (
`ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp',
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING
) WITH (
'connector' = 'kafka',
...
'format' = 'json',
'json.ignore-parse-errors' = 'true'
)
값 포맷은 다음 데이터 타입으로 구성됩니다:
ROW<`user_id` BIGINT, `item_id` BIGINT, `behavior` STRING>
키와 값 포맷 (Key and Value Format)
다음 예시는 키와 값 포맷의 지정·구성 방법을 보여줍니다. 포맷 옵션은 'key' 또는 'value'에 포맷 식별자를 더한 접두사로 시작합니다.
CREATE TABLE KafkaTable (
`ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp',
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING
) WITH (
'connector' = 'kafka',
...
'key.format' = 'json',
'key.json.ignore-parse-errors' = 'true',
'key.fields' = 'user_id;item_id',
'value.format' = 'json',
'value.json.fail-on-missing-field' = 'false',
'value.fields-include' = 'ALL'
)
키 포맷은 'key.fields'(';'를 구분자로 사용)에 나열된 필드를 같은 순서로 포함합니다. 따라서 다음 데이터 타입으로 구성됩니다:
ROW<`user_id` BIGINT, `item_id` BIGINT>
값 포맷이 'value.fields-include' = 'ALL'로 구성되었으므로 키 필드도 값 포맷의 데이터 타입에 들어갑니다:
ROW<`user_id` BIGINT, `item_id` BIGINT, `behavior` STRING>
겹치는 포맷 필드 (Overlapping Format Fields)
키와 값 포맷 둘 다 같은 이름의 필드를 포함하면 커넥터는 스키마 정보만으로 테이블 컬럼을 키·값 필드로 나눌 수 없습니다. 'key.fields-prefix' 옵션을 사용하면 키 포맷을 구성할 때 원래 이름을 유지하면서 테이블 스키마에서 키 컬럼에 고유한 이름을 부여할 수 있습니다.
다음 예시는 둘 다 version 필드를 포함하는 키와 값 포맷을 보여줍니다:
CREATE TABLE KafkaTable (
`k_version` INT,
`k_user_id` BIGINT,
`k_item_id` BIGINT,
`version` INT,
`behavior` STRING
) WITH (
'connector' = 'kafka',
...
'key.format' = 'json',
'key.fields-prefix' = 'k_',
'key.fields' = 'k_version;k_user_id;k_item_id',
'value.format' = 'json',
'value.fields-include' = 'EXCEPT_KEY'
)
값 포맷은 'EXCEPT_KEY' 모드로 구성되어야 합니다. 포맷은 다음 데이터 타입으로 구성됩니다:
key format:
ROW<`version` INT, `user_id` BIGINT, `item_id` BIGINT>
value format:
ROW<`version` INT, `behavior` STRING>
토픽과 파티션 발견 (Topic and Partition Discovery)
구성 옵션 topic과 topic-pattern은 소스에서 소비할 토픽 또는 토픽 패턴을 지정합니다. topic 구성 옵션은 'topic-1;topic-2'처럼 세미콜론 구분자를 사용한 토픽 목록을 받을 수 있습니다. topic-pattern 구성 옵션은 정규식을 사용해 일치하는 토픽을 발견합니다. 예를 들어 topic-pattern이 test-topic-[0-9]이면 지정된 정규식(test-topic-로 시작하고 한 자리 숫자로 끝나는)과 일치하는 이름의 모든 토픽이 작업 시작 시 컨슈머가 구독합니다.
작업 시작 후 동적으로 생성된 토픽을 컨슈머가 발견할 수 있게 하려면 scan.topic-partition-discovery.interval에 음수가 아닌 값을 설정하세요. 이렇게 하면 지정된 패턴과도 일치하는 새 토픽의 파티션을 컨슈머가 발견할 수 있습니다.
토픽과 파티션 발견에 대한 자세한 내용은 Kafka DataStream 커넥터 문서를 참고하세요.
토픽 목록과 토픽 패턴은 소스에서만 동작함에 유의하세요. 싱크에서 Flink는 현재 단일 토픽만 지원합니다.
시작 읽기 위치 (Start Reading Position)
구성 옵션 scan.startup.mode는 Kafka 컨슈머의 시작 모드를 지정합니다. 유효한 값은:
group-offsets: 특정 컨슈머 그룹의 ZK / Kafka 브로커에 커밋된 오프셋부터 시작.earliest-offset: 가능한 가장 이른 오프셋부터 시작.latest-offset: 가장 최근 오프셋부터 시작.timestamp: 각 파티션의 사용자 제공 타임스탬프부터 시작.specific-offsets: 각 파티션의 사용자 제공 특정 오프셋부터 시작.
기본 옵션 값은 group-offsets로, ZK / Kafka 브로커의 마지막 커밋 오프셋부터 소비함을 나타냅니다.
timestamp가 지정되면 1970년 1월 1일 00:00:00.000 GMT 이후의 밀리초 단위 특정 시작 타임스탬프를 지정하기 위해 또 다른 구성 옵션 scan.startup.timestamp-millis가 필요합니다.
specific-offsets가 지정되면 각 파티션의 특정 시작 오프셋을 지정하기 위해 또 다른 구성 옵션 scan.startup.specific-offsets가 필요합니다. 예를 들어 옵션 값 partition:0,offset:42;partition:1,offset:300은 파티션 0의 오프셋 42, 파티션 1의 오프셋 300을 나타냅니다.
유한 끝 위치 (Bounded Ending Position)
구성 옵션 scan.bounded.mode는 Kafka 컨슈머의 유한 모드를 지정합니다. 유효한 값은:
group-offsets: 특정 컨슈머 그룹의 ZooKeeper / Kafka 브로커에 커밋된 오프셋으로 유한. 주어진 파티션에서 소비 시작 시 평가됩니다.latest-offset: 최신 오프셋으로 유한. 주어진 파티션에서 소비 시작 시 평가됩니다.timestamp: 사용자 제공 타임스탬프로 유한.specific-offsets: 각 파티션의 사용자 제공 특정 오프셋으로 유한.
구성 옵션 값 scan.bounded.mode가 설정되지 않으면 기본값은 무한(unbounded) 테이블입니다.
timestamp가 지정되면 1970년 1월 1일 00:00:00.000 GMT 이후의 밀리초 단위 특정 유한 타임스탬프를 지정하기 위해 또 다른 구성 옵션 scan.bounded.timestamp-millis가 필요합니다.
specific-offsets가 지정되면 각 파티션의 특정 유한 오프셋을 지정하기 위해 또 다른 구성 옵션 scan.bounded.specific-offsets가 필요합니다. 예를 들어 옵션 값 partition:0,offset:42;partition:1,offset:300은 파티션 0의 오프셋 42, 파티션 1의 오프셋 300을 나타냅니다. 파티션에 대한 오프셋이 제공되지 않으면 해당 파티션에서 소비하지 않습니다.
CDC Changelog 소스 (CDC Changelog Source)
Flink는 Kafka를 CDC changelog 소스로 네이티브로 지원합니다. Kafka 토픽의 메시지가 CDC 도구를 사용해 다른 데이터베이스에서 캡처한 변경 이벤트라면, 해당 Flink CDC 포맷을 사용해 메시지를 Flink SQL 테이블에 대한 INSERT/UPDATE/DELETE 문으로 해석할 수 있습니다.
changelog 소스는 데이터베이스에서 다른 시스템으로 증분 데이터를 동기화하거나, 로그 감사, 데이터베이스의 구체화 뷰, 데이터베이스 테이블 변경 이력의 시간 조인 등 많은 경우에 매우 유용한 기능입니다.
Flink는 여러 CDC 포맷을 제공합니다:
- debezium
- canal
- maxwell
싱크 파티셔닝 (Sink Partitioning)
구성 옵션 sink.partitioner는 Flink의 파티션에서 Kafka의 파티션으로의 출력 파티셔닝을 지정합니다. 기본적으로 Flink는 레코드를 파티셔닝하기 위해 Kafka 기본 파티셔너를 사용합니다. null 키가 있는 레코드에는 sticky 파티션 전략을 사용하고, 키가 정의된 레코드에는 murmur2 해시로 파티션을 계산합니다.
행의 파티션 라우팅을 제어하기 위해 커스텀 싱크 파티셔너를 제공할 수 있습니다. fixed 파티셔너는 같은 Flink 파티션의 레코드를 같은 Kafka 파티션에 써서 네트워크 연결 비용을 줄일 수 있습니다.
일관성 보장 (Consistency guarantees)
기본적으로 Kafka 싱크는 체크포인팅이 활성화된 상태로 쿼리가 실행되면 at-least-once 보장으로 데이터를 Kafka 토픽에 넣습니다.
Flink의 체크포인팅이 활성화되면 kafka 커넥터는 exactly-once 전달 보장을 제공할 수 있습니다.
Flink의 체크포인팅 활성화 외에도 적절한 sink.delivery-guarantee 옵션을 전달해 선택한 세 가지 동작 모드 중 하나를 고를 수 있습니다:
none: Flink는 아무것도 보장하지 않습니다. 생성된 레코드는 유실되거나 중복될 수 있습니다.at-least-once(기본 설정): 레코드가 유실되지 않음을 보장합니다(중복될 수는 있음).exactly-once: exactly-once 의미론을 제공하기 위해 Kafka 트랜잭션을 사용합니다. 트랜잭션으로 Kafka에 쓸 때는 Kafka에서 레코드를 소비하는 모든 애플리케이션에 대해 (read_uncommitted또는read_committed중) 원하는isolation.level— 후자가 기본값 — 을 설정하는 것을 잊지 마세요.
전달 보장에 대한 더 많은 주의 사항은 Kafka 문서를 참고하세요.
소스 파티션별 워터마크 (Source Per-Partition Watermarks)
Flink는 Kafka에 대한 파티션별 워터마크 발행을 지원합니다. 워터마크는 Kafka 컨슈머 내부에서 생성됩니다. 파티션별 워터마크는 스트리밍 셔플 중 워터마크가 병합되는 것과 같은 방식으로 병합됩니다. 소스의 출력 워터마크는 읽는 파티션 중 최소 워터마크로 결정됩니다. 토픽의 일부 파티션이 유휴이면 워터마크 생성기가 진행되지 않습니다. 테이블 구성에서 'table.exec.source.idle-timeout' 옵션을 설정해 이 문제를 완화할 수 있습니다.
자세한 내용은 Kafka watermark strategies를 참고하세요.
보안 (Security)
암호화와 인증을 포함한 보안 구성을 활성화하려면 테이블 옵션에서 "properties." 접두사로 보안 구성을 설정하기만 하면 됩니다. 아래 코드 스니펫은 SQL client JAR을 사용할 때 SASL 메커니즘으로 PLAIN을 사용하고 JAAS 구성을 제공하도록 Kafka 테이블을 구성하는 방법을 보여줍니다:
CREATE TABLE KafkaTable (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING,
`ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp'
) WITH (
'connector' = 'kafka',
...
'properties.security.protocol' = 'SASL_PLAINTEXT',
'properties.sasl.mechanism' = 'PLAIN',
'properties.sasl.jaas.config' = 'org.apache.kafka.common.security.plain.PlainLoginModule required username="username" password="password";'
)
더 복잡한 예시로, SQL client JAR을 사용할 때 보안 프로토콜로 SASL_SSL을 사용하고 SASL 메커니즘으로 SCRAM-SHA-256을 사용합니다:
CREATE TABLE KafkaTable (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING,
`ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp'
) WITH (
'connector' = 'kafka',
...
'properties.security.protocol' = 'SASL_SSL',
/* SSL configurations */
/* Configure the path of truststore (CA) provided by the server */
'properties.ssl.truststore.location' = '/path/to/kafka.client.truststore.jks',
'properties.ssl.truststore.password' = 'test1234',
/* Configure the path of keystore (private key) if client authentication is required */
'properties.ssl.keystore.location' = '/path/to/kafka.client.keystore.jks',
'properties.ssl.keystore.password' = 'test1234',
/* SASL configurations */
/* Set SASL mechanism as SCRAM-SHA-256 */
'properties.sasl.mechanism' = 'SCRAM-SHA-256',
/* Set JAAS configurations */
'properties.sasl.jaas.config' = 'org.apache.kafka.common.security.scram.ScramLoginModule required username="username" password="password";'
)
Kafka 클라이언트 의존성을 리로케이션하면 sasl.jaas.config의 로그인 모듈 클래스 경로가 다를 수 있으므로, JAR의 모듈 실제 클래스 경로로 다시 작성해야 할 수 있음에 유의하세요. SQL client JAR은 Kafka 클라이언트 의존성을 org.apache.flink.kafka.shaded.org.apache.kafka로 리로케이션하므로, 위 코드 스니펫의 plain 로그인 모듈 경로는 SQL client JAR을 사용할 때 org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule이어야 합니다.
보안 구성에 대한 자세한 설명은 Apache Kafka 문서의 "Security" 섹션을 참고하세요.
데이터 타입 매핑 (Data Type Mapping)
Kafka는 메시지 키와 값을 바이트로 저장하므로 Kafka에는 스키마나 데이터 타입이 없습니다. Kafka 메시지는 csv, json, avro 같은 포맷으로 역직렬화·직렬화됩니다. 따라서 데이터 타입 매핑은 특정 포맷에 의해 결정됩니다. 자세한 내용은 Formats 페이지를 참고하세요.