Protobuf 포맷
Protobuf 포맷 (Protobuf Format)
Protocol Buffers Protobuf 포맷은 Protobuf 생성 클래스를 기반으로 Protobuf 데이터를 읽고 쓸 수 있게 해줍니다.
출처: 문서
본문
이 포맷은 직렬화 스키마(Serialization Schema) 와 역직렬화 스키마(Deserialization Schema) 로 사용할 수 있습니다.
Protocol Buffers Protobuf 포맷은 Protobuf 생성 클래스를 기반으로 Protobuf 데이터를 읽고 쓸 수 있게 해줍니다.
의존성
Protobuf 포맷을 사용하려면 빌드 자동화 도구(예: Maven 또는 SBT)를 사용하는 프로젝트와 SQL JAR 번들을 사용하는 SQL Client 모두에 다음 의존성이 필요합니다.
Maven 의존성:
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-protobuf</artifactId>
<version>2.3.0</version>
</dependency>
SQL Client: 다운로드
Protobuf 포맷으로 테이블 생성하기
Kafka 커넥터와 Protobuf 포맷을 사용해 테이블을 만드는 예는 다음과 같습니다.
아래는 proto 정의 파일입니다:
syntax = "proto2";
package com.example;
option java_package = "com.example";
option java_multiple_files = true;
message SimpleTest {
optional int64 uid = 1;
optional string name = 2;
optional int32 category_type = 3;
optional bytes content = 4;
optional double price = 5;
map<int64, InnerMessageTest> value_map = 6;
repeated InnerMessageTest value_arr = 7;
optional Corpus corpus_int = 8;
optional Corpus corpus_str = 9;
message InnerMessageTest{
optional int64 v1 =1;
optional int32 v2 =2;
}
enum Corpus {
UNIVERSAL = 0;
WEB = 1;
IMAGES = 2;
LOCAL = 3;
NEWS = 4;
PRODUCTS = 5;
VIDEO = 7;
}
}
protoc명령을 사용해.proto파일을 Java 클래스로 컴파일합니다.- 그런 다음 클래스를 컴파일하고 패키징합니다(proto-java 를 jar 에 패키징할 필요는 없습니다).
- 마지막으로 클래스패스에
jar를 제공해야 합니다. 예: sql-client 에서-j로 전달.
CREATE TABLE simple_test (
uid BIGINT,
name STRING,
category_type INT,
content BINARY,
price DOUBLE,
value_map map<BIGINT, row<v1 BIGINT, v2 INT>>,
value_arr array<row<v1 BIGINT, v2 INT>>,
corpus_int INT,
corpus_str STRING
) WITH (
'connector' = 'kafka',
'topic' = 'user_behavior',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'format' = 'protobuf',
'protobuf.message-class-name' = 'com.example.SimpleTest',
'protobuf.ignore-parse-errors' = 'true'
)
포맷 옵션 (Format Options)
| 옵션 | 필수 | 전달(Forwarded) | 기본값 | 유형 | 설명 |
|---|---|---|---|---|---|
format |
required | no | (none) | String | 사용할 포맷을 지정합니다. 여기서는 'protobuf' 여야 합니다. |
protobuf.message-class-name |
required | no | (none) | String | Protobuf 생성 클래스의 전체 이름입니다. 이름은 proto 정의 파일의 메시지 이름과 일치해야 합니다. 내부 클래스 이름에는 $ 가 지원됩니다(예: 'com.exmample.OuterClass$MessageClass'). |
protobuf.ignore-parse-errors |
optional | no | false | Boolean | 실패 대신 파싱 오류가 있는 행을 건너뛸지 여부를 나타내는 선택적 플래그입니다. |
protobuf.read-default-values |
optional | yes | false | Boolean | 이 값을 true 로 설정하면 포맷이 빈 값을 proto 파일에 정의된 기본값으로 읽습니다. false 로 설정하면 데이터 요소가 이진 protobuf 메시지에 없을 때 null 값을 생성합니다. Flink 의 현재 protobuf 버전(4.32.1)에서는 proto3 에 대해 필드 존재(field presence)가 제대로 지원되어 비기본 타입의 null 처리가 가능합니다. 이 값을 true 로 설정하면 스키마 복잡성과 메시지 크기에 따라 역직렬화 성능이 훨씬 느려질 수 있음을 유의하세요. |
protobuf.write-null-string-literal |
optional | no | "" | String | protobuf 데이터로 직렬화할 때, null 값이 있는 경우 Protobuf 의 array/map 에서 사용할 문자열 리터럴을 지정하는 선택적 구성입니다. |
데이터 타입 매핑 (Data Type Mapping)
아래 표는 Flink 타입에서 Protobuf 타입으로의 매핑을 나열합니다:
| Flink SQL 타입 | Protobuf 타입 | 설명 |
|---|---|---|
CHAR / VARCHAR / STRING |
string |
|
BOOLEAN |
bool |
|
BINARY / VARBINARY |
bytes |
|
INT |
int32 |
|
BIGINT |
int64 |
|
FLOAT |
float |
|
DOUBLE |
double |
|
ARRAY |
repeated |
요소는 null 일 수 없으며, 문자열 기본값은 write-null-string-literal 로 지정할 수 있습니다. |
MAP |
map |
키나 값은 null 일 수 없으며, 문자열 기본값은 write-null-string-literal 로 지정할 수 있습니다. |
ROW |
message |
|
VARCHAR / CHAR / TINYINT / SMALLINT / INTEGER / BIGINT |
enum |
protobuf 의 enum 값은 그에 따라 flink row 의 문자열 또는 숫자에 매핑할 수 있습니다. |
ROW<seconds BIGINT, nanos INT> |
google.protobuf.timestamp |
google.protobuf.timestamp 타입은 row 타입과 protobuf 정의를 사용해 UTC epoch 시간의 나노초 해상도로 초와 초의 분수를 매핑할 수 있습니다. |
Null 값
protobuf 는 map 과 array 에서 null 값을 허용하지 않으므로, Flink Rows 를 Protobuf 로 변환할 때 기본값을 자동 생성해야 합니다.
| Protobuf 데이터 타입 | 기본값 |
|---|---|
| int32 / int64 / float / double | 0 |
| string | "" |
| bool | false |
| enum | enum 의 첫 번째 요소 |
| binary | ByteString.EMPTY |
| message | MESSAGE.getDefaultInstance() |
OneOf 필드
직렬화 과정에서 같은 one-of 그룹의 Flink 필드가 최대 하나의 유효한 값만 포함한다는 보장은 없습니다. 직렬화 시 각 필드는 Flink 스키마 순서대로 설정되므로, 같은 one-of 그룹에서 더 높은 위치의 필드가 더 낮은 위치의 필드를 덮어씁니다.
지원되는 Protobuf 버전
Flink 는 protobuf-java 4.32.1(Protocol Buffers 버전 32에 해당)을 사용하며, 다음을 지원합니다:
- Proto2 및 Proto3 문법: 전통적인
syntax = "proto2"및syntax = "proto3"정의 - Protobuf Editions: Protocol Buffers v27+ 에서 도입된 새로운
edition = "2023"및edition = "2024"문법 - 개선된 proto3 필드 존재 감지: 이전 protobuf 버전의 제한 없이 optional 필드를 더 잘 처리
Protobuf Editions 사용하기
Protobuf Editions 는 proto2 와 proto3 기능을 결합한 통합 문법을 제공합니다. .proto 파일에서 Editions 를 사용하는 경우 Flink 는 이를 완전히 지원합니다:
edition = "2023";
package com.example;
option java_package = "com.example";
option java_multiple_files = true;
message SimpleTest {
int64 uid = 1;
string name = 2 [features.field_presence = EXPLICIT];
// ... rest of your message definition
}
Editions 는 proto2 및 proto3 와의 하위 호환성을 유지하면서 파일, 메시지 또는 필드 수준에서 기능 동작을 세밀하게 제어할 수 있게 해줍니다. 자세한 내용은 Protobuf Editions 문서 를 참고하세요.
추가 리소스
Protocol Buffers 에 대한 자세한 내용은 다음을 참고하세요:
- Language Guide (proto2)
- Language Guide (proto3)
- Language Guide (Editions) - 새 Editions 문법용
- Protobuf Editions Overview - Editions 의 동기와 이점 이해