JSON 포맷
JSON 포맷 (Json format)
JSON 포맷은 Flink 작업에서 JSON 레코드를 읽고 쓰는 데 사용됩니다. JsonSerializationSchema/JsonDeserializationSchema를 통해 JSON 직렬화/역직렬화를 지원하며, 내부적으로 Jackson 라이브러리를 사용합니다. POJO와 ObjectNode를 포함해 Jackson이 지원하는 모든 타입을 처리할 수 있습니다.
출처: 문서
본문
JSON 포맷을 사용하려면 프로젝트에 Flink JSON 의존성을 추가해야 합니다.
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-json</artifactId>
<version>2.3.0</version>
<scope>provided</scope>
</dependency>
PyFlink 사용자는 작업에서 이를 직접 사용할 수 있습니다.
Flink는 JsonSerializationSchema/JsonDeserializationSchema를 통해 JSON 레코드 읽기/쓰기를 지원합니다. 이들은 Jackson 라이브러리를 활용하며, POJO와 ObjectNode를 포함해(그리고 그에 국한되지 않고) Jackson이 지원하는 모든 타입을 지원합니다.
JsonDeserializationSchema는 DeserializationSchema를 지원하는 모든 커넥터와 함께 사용할 수 있습니다.
예를 들어 KafkaSource와 함께 사용해 POJO를 역직렬화하는 방법은 다음과 같습니다.
JsonDeserializationSchema<SomePojo> jsonFormat=new JsonDeserializationSchema<>(SomePojo.class);
KafkaSource<SomePojo> source=
KafkaSource.<SomePojo>builder()
.setValueOnlyDeserializer(jsonFormat)
...
JsonSerializationSchema는 SerializationSchema를 지원하는 모든 커넥터와 함께 사용할 수 있습니다.
예를 들어 KafkaSink와 함께 사용해 POJO를 직렬화하는 방법은 다음과 같습니다.
JsonSerializationSchema<SomePojo> jsonFormat=new JsonSerializationSchema<>();
KafkaSink<SomePojo> source =
KafkaSink.<SomePojo>builder()
.setRecordSerializer(
new KafkaRecordSerializationSchemaBuilder<>()
.setValueSerializationSchema(jsonFormat)
...
사용자 정의 Mapper (Custom Mapper)
두 스키마 모두 SerializableSupplier<ObjectMapper>를 받는 생성자가 있으며, 이는 object mapper의 팩토리 역할을 합니다. 이 팩토리로 생성된 mapper를 완전히 제어할 수 있고, 다양한 Jackson 기능을 활성화/비활성화하거나 모듈을 등록해 지원되는 타입의 집합을 확장하거나 추가 기능을 더할 수 있습니다.
JsonSerializationSchema<SomeClass> jsonFormat=new JsonSerializationSchema<>(
() -> new ObjectMapper()
.enable(SerializationFeature.ORDER_MAP_ENTRIES_BY_KEYS))
.registerModule(new ParameterNamesModule());
Python
PyFlink에서 JsonRowSerializationSchema와 JsonRowDeserializationSchema는 Row 타입에 대한 내장 지원입니다. 예를 들어 KafkaSource와 KafkaSink에서 사용하는 방법은 다음과 같습니다.
row_type_info = Types.ROW_NAMED(['name', 'age'], [Types.STRING(), Types.INT()])
json_format = JsonRowDeserializationSchema.builder().type_info(row_type_info).build()
source = KafkaSource.builder() \
.set_value_only_deserializer(json_format) \
.build()
row_type_info = Types.ROW_NAMED(['name', 'age'], [Types.STRING(), Types.INT()])
json_format = JsonRowSerializationSchema.builder().with_type_info(row_type_info).build()
sink = KafkaSink.builder() \
.set_record_serializer(
KafkaRecordSerializationSchema.builder()
.set_topic('test')
.set_value_serialization_schema(json_format)
.build()
) \
.build()