1. JsonDeserializationSchema:KafkaSource 中反序列化 POJO

JsonDeserializationSchema 实现了 Flink 的 DeserializationSchema,因此只要某个 connector 支持 DeserializationSchema,你就能直接使用它。

典型用法:KafkaSource 只消费 value,反序列化成 POJO:

JsonDeserializationSchema<SomePojo> jsonFormat =
    new JsonDeserializationSchema<>(SomePojo.class);

KafkaSource<SomePojo> source =
    KafkaSource.<SomePojo>builder()
        .setValueOnlyDeserializer(jsonFormat)
        // ...
        .build();

适用场景:

  • Kafka 的 value 是 JSON
  • 你希望在 DataStream 里直接拿到业务对象 SomePojo

工程建议:

  • POJO 字段尽量使用包装类型(Integer/Long)应对字段缺失或 null
  • 为了兼容字段变动,可以配合 ObjectMapper 设置忽略未知字段(见第 3 节)

2. JsonSerializationSchema:KafkaSink 中序列化 POJO

写回 Kafka 时,JsonSerializationSchema 实现了 SerializationSchema,可用于任何支持 SerializationSchema 的 connector。

典型用法:KafkaSink 写 value,序列化 POJO 为 JSON:

JsonSerializationSchema<SomePojo> jsonFormat =
    new JsonSerializationSchema<>();

KafkaSink<SomePojo> sink =
    KafkaSink.<SomePojo>builder()
        .setRecordSerializer(
            new KafkaRecordSerializationSchemaBuilder<SomePojo>()
                .setValueSerializationSchema(jsonFormat)
                // ...
                .build()
        )
        .build();

适用场景:

  • 你希望下游系统继续消费 JSON
  • 你不想自己手写 Jackson 序列化逻辑

3. 自定义 ObjectMapper:控制 Jackson 行为(非常常用)

Flink 允许你通过构造函数传入 SerializableSupplier<ObjectMapper> 来定制 mapper,相当于提供一个“ObjectMapper 工厂”。

你可以用它做很多工程级增强,比如:

  • 忽略未知字段(兼容上游 schema 变更)
  • 注册模块(Java 时间类型、参数名模块等)
  • 开启/关闭某些序列化特性(字段排序、空值处理等)

示例:自定义序列化 mapper,让 map key 有序,并注册模块:

JsonSerializationSchema<SomeClass> jsonFormat =
    new JsonSerializationSchema<>(
        () -> new ObjectMapper()
            .enable(SerializationFeature.ORDER_MAP_ENTRIES_BY_KEYS)
            .registerModule(new ParameterNamesModule())
    );

你也可以把“兼容字段变更”的设置加进去(强烈建议生产开启类似配置):

  • FAIL_ON_UNKNOWN_PROPERTIES 关闭
  • JavaTimeModule 等

(这里不展开写完整 mapper 配置,你只要知道:用 supplier 你就能完全掌控 Jackson。)

4. PyFlink:Row 类型用 JsonRowSerializationSchema / JsonRowDeserializationSchema

在 PyFlink 中,Flink 内置了 Row 的 JSON Schema:

  • JsonRowDeserializationSchema
  • JsonRowSerializationSchema

这对 Python 流处理特别友好,因为 Python 侧更常操作 Row 而不是 POJO 类。

KafkaSource:JSON -> Row

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

KafkaSink:Row -> JSON

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

适用场景:

  • Python 处理流数据,行结构清晰
  • Kafka 中 value 为 JSON

5. 选型建议:POJO vs ObjectNode vs Row

  • Java POJO:类型安全、IDE 友好、适合稳定 schema 的业务流
  • ObjectNode:更灵活,适合 schema 频繁变化、半结构化数据
  • PyFlink Row:Python 生态更顺手,适合表/行式处理
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐