JSON 포맷

JSON 포맷 (Json format)

JSON 포맷은 Flink 작업에서 JSON 레코드를 읽고 쓰는 데 사용됩니다. JsonSerializationSchema/JsonDeserializationSchema를 통해 JSON 직렬화/역직렬화를 지원하며, 내부적으로 Jackson 라이브러리를 사용합니다. POJOObjectNode를 포함해 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 라이브러리를 활용하며, POJOObjectNode를 포함해(그리고 그에 국한되지 않고) Jackson이 지원하는 모든 타입을 지원합니다.

JsonDeserializationSchemaDeserializationSchema를 지원하는 모든 커넥터와 함께 사용할 수 있습니다.

예를 들어 KafkaSource와 함께 사용해 POJO를 역직렬화하는 방법은 다음과 같습니다.

JsonDeserializationSchema<SomePojo> jsonFormat=new JsonDeserializationSchema<>(SomePojo.class);
KafkaSource<SomePojo> source=
    KafkaSource.<SomePojo>builder()
        .setValueOnlyDeserializer(jsonFormat)
        ...

JsonSerializationSchemaSerializationSchema를 지원하는 모든 커넥터와 함께 사용할 수 있습니다.

예를 들어 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에서 JsonRowSerializationSchemaJsonRowDeserializationSchemaRow 타입에 대한 내장 지원입니다. 예를 들어 KafkaSourceKafkaSink에서 사용하는 방법은 다음과 같습니다.

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()

더 알아보기 (Learn more)