SerDe 사용하기
SerDe 사용하기
Pulsar 함수는 Pulsar 토픽에 데이터를 게시하거나 소비할 때 SerDe(직렬화·역직렬화)를 사용해요. SerDe가 기본적으로 동작하는 방식은 함수를 작성한 언어(Java 또는 Python)에 따라 달라요. 다만 두 언어 모두에서 더 복잡한 애플리케이션 전용 타입을 처리하려면 나만의 SerDe 로직을 작성할 수 있어요.
SerDe를 이해하면 함수 사이에서 객체를 그대로 주고받을 수 있어요. 예를 들어 트윗(Tweet) 객체를 함수 간에 직접 전달하고 싶다면 커스텀 SerDe 클래스를 하나 만들어 주면 된답니다.
출처: 문서
본문
Java 함수용 SerDe (Use SerDe for Java functions)
Java 함수에서 기본으로 지원되는 내장 타입은 string, double, integer, float, long, short, byte예요.
Java 타입을 커스터마이즈하려면 다음 인터페이스를 구현해야 해요.
public interface SerDe<T> {
T deserialize(byte[] input);
byte[] serialize(T input);
}
Java 함수에서 SerDe는 다음과 같이 동작해요.
- 입력·출력 토픽에 스키마가 있으면 Pulsar 함수는 그 스키마를 SerDe에 사용해요.
- 입력·출력 토픽이 없으면 Pulsar 함수는 다음 규칙으로 SerDe를 결정해요.
- 스키마 타입이 지정되어 있으면, Pulsar 함수는 지정된 스키마 타입을 사용해요.
- SerDe가 지정되어 있으면, Pulsar 함수는 지정된 SerDe를 사용하고 입력·출력 토픽의 스키마 타입은
byte가 돼요. - 스키마 타입도 SerDe도 지정되지 않으면, Pulsar 함수는 내장 SerDe를 사용해요. 비원시(non-primitive) 스키마 타입의 경우 내장 SerDe가 객체를 JSON 형식으로 직렬화·역직렬화해요.
예를 들어 트윗 객체를 처리하는 함수를 작성한다고 생각해 볼게요. Java의 Tweet 클래스 예시는 다음과 같아요.
public class Tweet {
private String username;
private String tweetContent;
public Tweet(String username, String tweetContent) {
this.username = username;
this.tweetContent = tweetContent;
}
// Standard setters and getters
}
트윗 객체를 함수 사이에서 직접 전달하려면 커스텀 SerDe 클래스가 필요해요. 아래 예시에서 Tweet 객체는 기본적으로 문자열이고, username과 tweet content는 |로 구분돼요.
package com.example.serde;
import org.apache.pulsar.functions.api.SerDe;
import java.util.regex.Pattern;
public class TweetSerde implements SerDe<Tweet> {
public Tweet deserialize(byte[] input) {
String s = new String(input);
String[] fields = s.split(Pattern.quote("|"));
return new Tweet(fields[0], fields[1]);
}
public byte[] serialize(Tweet input) {
return "%s|%s".format(input.getUsername(), input.getTweetContent()).getBytes();
}
}
특정 함수에 커스터마이즈된 SerDe를 적용하려면 다음을 해야 해요.
Tweet과TweetSerde클래스를 JAR로 패키징한다.- 함수를 배포할 때 JAR 경로와 SerDe 클래스 이름을 지정한다.
create 명령으로 커스텀 SerDe를 적용해 함수를 배포하는 예시예요.
bin/pulsar-admin functions create \
--jar /path/to/your.jar \
--output-serde-classname com.example.serde.TweetSerde \
# Other function attributes
알아두기: 커스텀 SerDe 클래스는 함수 JAR와 함께 패키징되어야 해요.
Python 함수용 SerDe (Use SerDe for Python functions)
Python에서 기본 SerDe는 identity(항등)예요. 즉 타입이 함수가 반환하는 타입 그대로 직렬화돼요.
예를 들어 클러스터 모드에서 함수를 배포할 때 SerDe를 다음과 같이 지정할 수 있어요.
bin/pulsar-admin functions create \
--tenant public \
--namespace default \
--name my_function \
--py my_function.py \
--classname my_function.MyFunction \
--custom-serde-inputs '{"input-topic-1":"Serde1","input-topic-2":"Serde2"}' \
--output-serde-classname Serde3 \
--output output-topic-1
이 경우 두 개의 입력 토픽(input-topic-1, input-topic-2)이 각각 서로 다른 SerDe 클래스에 매핑돼요(매핑은 JSON 문자열로 지정해야 해요). 출력 토픽 output-topic-1은 Serde3 클래스를 SerDe에 사용해요.
알아두기: 처리 로직과 SerDe 클래스를 포함한 모든 함수 관련 로직은 단일 Python 파일 안에 담겨야 해요.
다음 표는 Python 함수의 세 가지 SerDe 옵션을 정리한 거예요.
| SerDe 옵션 | 설명 | 사용 사례 |
|---|---|---|
| IdentitySerde (기본값) | 데이터를 바꾸지 않는 IdentitySerde를 사용해요. SerDe를 명시하지 않고 함수를 만들거나 실행하면 이 옵션이 사용돼요. | 문자열, 부울, 정수 같은 단순 타입을 다룰 때 |
| PickleSerDe | Python pickle을 SerDe에 사용하는 PickleSerDe를 사용해요. | 복잡하고 애플리케이션 전용인 타입을 다룰 때, pickle의 "최선 노력(best-effort)" 방식이 괜찮을 때 |
| Custom SerDe | 두 개의 메서드만 가진 기본 SerDe 클래스를 구현해 커스텀 SerDe 클래스를 만들어요 — serialize(객체를 바이트로), deserialize(바이트를 애플리케이션 전용 타입 객체로). |
성능이나 데이터 호환성 때문에 SerDe를 명시적으로 제어하고 싶을 때 |
예를 들어 트윗 객체를 처리하는 함수를 작성한다고 생각해 볼게요. Python의 Tweet 클래스 예시는 다음과 같아요.
class Tweet(object):
def __init__(self, username, tweet_content):
self.username = username
self.tweet_content = tweet_content
이 클래스를 Pulsar 함수에서 사용하는 방법은 두 가지예요.
- SerDe에 pickle 라이브러리를 적용하는
PickleSerDe를 지정한다. - 나만의 SerDe 클래스를 만든다. 다음은 그 예시예요.
from pulsar import SerDe
class TweetSerDe(SerDe):
def serialize(self, input):
return bytes("{0}|{1}".format(input.username, input.tweet_content))
def deserialize(self, input_bytes):
tweet_components = str(input_bytes).split('|')
return Tweet(tweet_components[0], tweet_componentsp[1])
더 자세한 내용은 코드 예시를 참고해요.
더 알아보기 (Learn more)
- 함수의 직렬화에 영향을 주는 스키마 개념은 스키마 레지스트리 문서를 참고해요.
- Java 함수에 쓰는 스키마 타입은 Java 함수 문서에서 확인할 수 있어요.
- Python 함수의 SerDe 상세는 Python 함수 문서를 살펴보세요.